From 4e0eff36536ddc365c9d55730d6063cff417bb50 Mon Sep 17 00:00:00 2001 From: zhangtao <9480807882@qq.com> Date: Mon, 23 Feb 2026 06:43:38 +0800 Subject: [PATCH] =?UTF-8?q?feat(workflow):=20=E9=87=8D=E6=9E=84=E5=B7=A5?= =?UTF-8?q?=E4=BD=9C=E6=B5=81=E6=A8=A1=E5=9D=97=E5=B9=B6=E7=A7=BB=E9=99=A4?= =?UTF-8?q?=E9=85=8D=E7=BD=AE=E8=A1=A8=E5=8D=95=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit refactor(node): 移除节点模型中的config_schema字段 refactor(workflow): 简化工作流模型和CRUD操作 feat(composables): 重构并增强任务相关组合式函数 style(settings): 优化设置抽屉样式 fix(auth): 修复配置加载逻辑 chore: 清理无用文件和代码 --- backend/app/core/ap_scheduler.py | 8 +- .../app/plugin/module_task/job/controller.py | 2 +- backend/app/plugin/module_task/job/schema.py | 3 +- backend/app/plugin/module_task/node/model.py | 3 +- backend/app/plugin/module_task/node/schema.py | 1 - .../app/plugin/module_task/node/service.py | 3 +- .../plugin/module_task/workflow/controller.py | 489 +---- .../app/plugin/module_task/workflow/crud.py | 129 +- .../app/plugin/module_task/workflow/engine.py | 397 ++++ .../app/plugin/module_task/workflow/model.py | 85 +- .../app/plugin/module_task/workflow/schema.py | 172 +- .../plugin/module_task/workflow/service.py | 904 +------- backend/app/scripts/data/sys_menu.json | 278 +-- backend/app/utils/upload_util.py | 4 - .../data_level0.bin | Bin 628400 -> 0 bytes .../header.bin | Bin 100 -> 0 bytes .../length.bin | Bin 400 -> 0 bytes .../link_lists.bin | 0 backend/data/chroma/chroma.sqlite3 | Bin 208896 -> 0 bytes .../mysql/fastapiadmin_2026-01-03_201546.sql | 894 -------- .../mysql/fastapiadmin_2026-02-23_063912.sql | 963 ++++++++ ...sql => fastapiadmin_2026-02-23_064207.sql} | 1941 ++++++++++------- frontend/src/api/module_task/node.ts | 6 +- frontend/src/api/module_task/workflow.ts | 291 +-- frontend/src/composables/index.ts | 6 + .../src/composables/{ => task}/useNodeDrag.ts | 10 +- .../{ => task}/useNodeOperations.ts | 0 .../composables/{ => task}/usePerformance.ts | 0 .../{ => task}/useWorkflowHistory.ts | 0 .../src/layouts/components/Settings/index.vue | 7 +- frontend/src/store/modules/config.store.ts | 28 +- .../src/views/module_system/auth/index.vue | 3 - frontend/src/views/module_task/node/index.vue | 4 +- .../workflow/components/DynamicNode.vue | 42 +- .../workflow/components/NodeConfigPanel.vue | 103 +- .../components/WorkflowDesignDrawer.vue | 95 +- .../workflow/components/WorkflowRunDrawer.vue | 903 -------- .../src/views/module_task/workflow/index.vue | 143 +- 38 files changed, 2810 insertions(+), 5107 deletions(-) create mode 100644 backend/app/plugin/module_task/workflow/engine.py delete mode 100644 backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/data_level0.bin delete mode 100644 backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/header.bin delete mode 100644 backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/length.bin delete mode 100644 backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/link_lists.bin delete mode 100644 backend/data/chroma/chroma.sqlite3 delete mode 100644 backend/sql/mysql/fastapiadmin_2026-01-03_201546.sql create mode 100644 backend/sql/mysql/fastapiadmin_2026-02-23_063912.sql rename backend/sql/postgres/{fastapiadmin_2026-01-03_213232.sql => fastapiadmin_2026-02-23_064207.sql} (63%) rename frontend/src/composables/{ => task}/useNodeDrag.ts (90%) rename frontend/src/composables/{ => task}/useNodeOperations.ts (100%) rename frontend/src/composables/{ => task}/usePerformance.ts (100%) rename frontend/src/composables/{ => task}/useWorkflowHistory.ts (100%) delete mode 100644 frontend/src/views/module_task/workflow/components/WorkflowRunDrawer.vue diff --git a/backend/app/core/ap_scheduler.py b/backend/app/core/ap_scheduler.py index 965038a0..0c6f4d20 100644 --- a/backend/app/core/ap_scheduler.py +++ b/backend/app/core/ap_scheduler.py @@ -2,8 +2,6 @@ import json from datetime import datetime from typing import Any -from redis.asyncio import Redis - from apscheduler.events import ( EVENT_ALL, EVENT_JOB_ERROR, @@ -25,15 +23,13 @@ from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.date import DateTrigger from apscheduler.triggers.interval import IntervalTrigger +from redis.asyncio import Redis -from app.common.enums import RedisInitKeyConfig from app.config.setting import settings from app.core.database import engine from app.core.logger import log -from app.core.redis_crud import RedisCURD -from app.utils.cron_util import CronUtil from app.plugin.module_task.node.model import NodeModel - +from app.utils.cron_util import CronUtil scheduler = AsyncIOScheduler() scheduler.configure( diff --git a/backend/app/plugin/module_task/job/controller.py b/backend/app/plugin/module_task/job/controller.py index c6d3182a..ecc34c1f 100644 --- a/backend/app/plugin/module_task/job/controller.py +++ b/backend/app/plugin/module_task/job/controller.py @@ -6,11 +6,11 @@ from fastapi.responses import JSONResponse from app.api.v1.module_system.auth.schema import AuthSchema from app.common.request import PaginationService from app.common.response import ResponseSchema, SuccessResponse +from app.core.ap_scheduler import SchedulerUtil from app.core.base_params import PaginationQueryParam from app.core.dependencies import AuthPermission from app.core.logger import log from app.core.router_class import OperationLogRoute -from app.core.ap_scheduler import SchedulerUtil from .schema import JobOutSchema, JobQueryParam from .service import JobService diff --git a/backend/app/plugin/module_task/job/schema.py b/backend/app/plugin/module_task/job/schema.py index a9da4f00..0ae817f5 100644 --- a/backend/app/plugin/module_task/job/schema.py +++ b/backend/app/plugin/module_task/job/schema.py @@ -34,12 +34,11 @@ class JobUpdateSchema(BaseModel): class JobOutSchema(JobCreateSchema, BaseSchema): """执行日志响应模型""" - + model_config = ConfigDict(from_attributes=True) ... - class JobQueryParam: """执行日志查询参数""" diff --git a/backend/app/plugin/module_task/node/model.py b/backend/app/plugin/module_task/node/model.py index d54e18a5..33895864 100644 --- a/backend/app/plugin/module_task/node/model.py +++ b/backend/app/plugin/module_task/node/model.py @@ -1,6 +1,6 @@ import enum -from sqlalchemy import Boolean, Integer, String, Text, JSON +from sqlalchemy import Boolean, Integer, String, Text from sqlalchemy.orm import Mapped, mapped_column from app.core.base_model import ModelMixin, UserMixin @@ -27,7 +27,6 @@ class NodeModel(ModelMixin, UserMixin): name: Mapped[str] = mapped_column(String(64), nullable=False, comment="节点名称") code: Mapped[str] = mapped_column(String(32), nullable=False, unique=True, comment="节点编码") category: Mapped[str] = mapped_column(String(32), default=NodeCategoryEnum.ACTION.value, comment="节点分类") - config_schema: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict, comment="配置表单Schema(JSON Schema)") jobstore: Mapped[str | None] = mapped_column(String(64), nullable=True, default="default", comment="存储器") executor: Mapped[str | None] = mapped_column(String(64), nullable=True, default="default", comment="执行器") trigger: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="触发器") diff --git a/backend/app/plugin/module_task/node/schema.py b/backend/app/plugin/module_task/node/schema.py index 166ffc91..da6c3b64 100644 --- a/backend/app/plugin/module_task/node/schema.py +++ b/backend/app/plugin/module_task/node/schema.py @@ -29,7 +29,6 @@ class NodeCreateSchema(BaseModel): end_date: str | None = Field(default=None, description="结束时间") code: str | None = Field(default=None, description="节点编码") category: str | None = Field(default=None, description="节点分类") - config_schema: dict | None = Field(default=None, description="节点配置") @model_validator(mode="after") def _validate_func(self): diff --git a/backend/app/plugin/module_task/node/service.py b/backend/app/plugin/module_task/node/service.py index 79a084f5..dae62224 100644 --- a/backend/app/plugin/module_task/node/service.py +++ b/backend/app/plugin/module_task/node/service.py @@ -1,7 +1,7 @@ from app.api.v1.module_system.auth.schema import AuthSchema +from app.core.ap_scheduler import SchedulerUtil from app.core.exceptions import CustomException from app.utils.cron_util import CronUtil -from app.core.ap_scheduler import SchedulerUtil from .crud import NodeCRUD from .schema import ( @@ -41,7 +41,6 @@ class NodeService: "name": obj.name, "code": obj.code, "category": obj.category, - "config_schema": obj.config_schema if obj.config_schema else {"fields": []}, "func": obj.func, "args": obj.args, "kwargs": obj.kwargs, diff --git a/backend/app/plugin/module_task/workflow/controller.py b/backend/app/plugin/module_task/workflow/controller.py index dcec63dc..f8729f22 100644 --- a/backend/app/plugin/module_task/workflow/controller.py +++ b/backend/app/plugin/module_task/workflow/controller.py @@ -17,17 +17,9 @@ from .schema import ( WorkflowOutSchema, WorkflowPublishSchema, WorkflowQueryParam, - WorkflowRunCreateSchema, - WorkflowRunLogOutSchema, - WorkflowRunLogQueryParam, - WorkflowRunOutSchema, - WorkflowRunQueryParam, - WorkflowRunUpdateSchema, WorkflowUpdateSchema, - WorkflowValidateResultSchema, - WorkflowValidateSchema, ) -from .service import WorkflowRunLogService, WorkflowRunService, WorkflowService +from .service import WorkflowService WorkflowRouter = APIRouter(route_class=OperationLogRoute, prefix="/workflow", tags=["工作流模块"]) @@ -194,104 +186,6 @@ async def publish_obj_controller( return SuccessResponse(data=result_dict, msg="发布工作流成功") -@WorkflowRouter.post( - "/validate", - summary="验证工作流", - description="验证工作流", - response_model=ResponseSchema[WorkflowValidateResultSchema], -) -async def validate_obj_controller( - data: WorkflowValidateSchema, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:validate"]))], -) -> JSONResponse: - """ - 验证工作流 - - 参数: - - data (WorkflowValidateSchema): 验证模型 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含验证结果的JSON响应 - """ - result_dict = await WorkflowService.validate_service(auth=auth, data=data) - log.info("验证工作流成功") - return SuccessResponse(data=result_dict.model_dump(), msg="验证工作流成功") - - -@WorkflowRouter.get( - "/templates", - summary="获取模板列表", - description="获取模板列表", - response_model=ResponseSchema[list[WorkflowOutSchema]], -) -async def get_templates_controller( - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:query"]))], -) -> JSONResponse: - """ - 获取模板列表 - - 参数: - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含模板列表的JSON响应 - """ - result_dict = await WorkflowService.templates_service(auth=auth) - log.info("获取模板列表成功") - return SuccessResponse(data=result_dict, msg="获取模板列表成功") - - -@WorkflowRouter.get( - "/export/{id}", - summary="导出工作流", - description="导出工作流", - response_model=ResponseSchema[dict], -) -async def export_obj_controller( - id: Annotated[int, Path(description="工作流ID")], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:export"]))], -) -> JSONResponse: - """ - 导出工作流 - - 参数: - - id (int): 工作流ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含导出工作流数据的JSON响应 - """ - result_dict = await WorkflowService.export_service(auth=auth, id=id) - log.info(f"导出工作流成功: {id}") - return SuccessResponse(data=result_dict, msg="导出工作流成功") - - -@WorkflowRouter.post( - "/import", - summary="导入工作流", - description="导入工作流", - response_model=ResponseSchema[WorkflowOutSchema], -) -async def import_obj_controller( - data: dict, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:import"]))], -) -> JSONResponse: - """ - 导入工作流 - - 参数: - - data (dict): 工作流数据 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含导入工作流详情的JSON响应 - """ - result_dict = await WorkflowService.import_service(auth=auth, data=data) - log.info(f"导入工作流成功: {result_dict.get('name')}") - return SuccessResponse(data=result_dict, msg="导入工作流成功") - - @WorkflowRouter.post( "/execute", summary="执行工作流", @@ -315,384 +209,3 @@ async def execute_obj_controller( result_dict = await WorkflowService.execute_service(auth=auth, data=data) log.info(f"执行工作流成功: {result_dict.get('workflow_id')}") return SuccessResponse(data=result_dict, msg="执行工作流成功") - - -@WorkflowRouter.post( - "/pause/{id}", - summary="暂停工作流", - description="暂停工作流", - response_model=ResponseSchema[None], -) -async def pause_obj_controller( - id: Annotated[int, Path(description="工作流ID")], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:pause"]))], -) -> JSONResponse: - """ - 暂停工作流 - - 参数: - - id (int): 工作流ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含暂停结果的JSON响应 - """ - await WorkflowService.pause_service(auth=auth, id=id) - log.info(f"暂停工作流成功: {id}") - return SuccessResponse(msg="暂停工作流成功") - - -@WorkflowRouter.post( - "/resume/{id}", - summary="恢复工作流", - description="恢复工作流", - response_model=ResponseSchema[None], -) -async def resume_obj_controller( - id: Annotated[int, Path(description="工作流ID")], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:resume"]))], -) -> JSONResponse: - """ - 恢复工作流 - - 参数: - - id (int): 工作流ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含恢复结果的JSON响应 - """ - await WorkflowService.resume_service(auth=auth, id=id) - log.info(f"恢复工作流成功: {id}") - return SuccessResponse(msg="恢复工作流成功") - - -@WorkflowRouter.post( - "/terminate/{id}", - summary="终止工作流", - description="终止工作流", - response_model=ResponseSchema[None], -) -async def terminate_obj_controller( - id: Annotated[int, Path(description="工作流ID")], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:terminate"]))], -) -> JSONResponse: - """ - 终止工作流 - - 参数: - - id (int): 工作流ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含终止结果的JSON响应 - """ - await WorkflowService.terminate_service(auth=auth, id=id) - log.info(f"终止工作流成功: {id}") - return SuccessResponse(msg="终止工作流成功") - - -@WorkflowRouter.get( - "/run/list", - summary="查询工作流运行记录列表", - description="查询工作流运行记录列表", - response_model=ResponseSchema[list[WorkflowRunOutSchema]], -) -async def get_workflow_run_list_controller( - page: Annotated[PaginationQueryParam, Depends()], - search: Annotated[WorkflowRunQueryParam, Depends()], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:query"]))], -) -> JSONResponse: - """ - 查询工作流运行记录列表 - - 参数: - - page (PaginationQueryParam): 分页查询参数 - - search (WorkflowRunQueryParam): 查询参数 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含工作流运行记录列表的JSON响应 - """ - result_dict = await WorkflowRunService.page_service( - auth=auth, - page_no=page.page_no, - page_size=page.page_size, - search=search.__dict__ if search else None, - order_by=page.order_by, - ) - return SuccessResponse(data=result_dict, msg="查询工作流运行记录列表成功") - - -@WorkflowRouter.get( - "/run/detail/{id}", - summary="查询工作流运行记录详情", - description="查询工作流运行记录详情", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def get_workflow_run_detail_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:query"]))], -) -> JSONResponse: - """ - 查询工作流运行记录详情 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含工作流运行记录详情的JSON响应 - """ - result_dict = await WorkflowRunService.detail_service(auth=auth, id=id) - return SuccessResponse(data=result_dict, msg="查询工作流运行记录详情成功") - - -@WorkflowRouter.post( - "/run/create", - summary="创建工作流运行记录", - description="创建工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def create_workflow_run_controller( - data: WorkflowRunCreateSchema, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:create"]))], -) -> JSONResponse: - """ - 创建工作流运行记录 - - 参数: - - data (WorkflowRunCreateSchema): 创建模型 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含创建的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.create_service(auth=auth, data=data) - return SuccessResponse(data=result_dict, msg="创建工作流运行记录成功") - - -@WorkflowRouter.put( - "/run/update/{id}", - summary="更新工作流运行记录", - description="更新工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def update_workflow_run_controller( - id: int, - data: WorkflowRunUpdateSchema, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], -) -> JSONResponse: - """ - 更新工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - data (WorkflowRunUpdateSchema): 更新模型 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含更新的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.update_service(auth=auth, id=id, data=data) - return SuccessResponse(data=result_dict, msg="更新工作流运行记录成功") - - -@WorkflowRouter.delete( - "/run/delete", - summary="删除工作流运行记录", - description="删除工作流运行记录", - response_model=ResponseSchema[None], -) -async def delete_workflow_run_controller( - ids: list[int], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:delete"]))], -) -> JSONResponse: - """ - 删除工作流运行记录 - - 参数: - - ids (list[int]): 工作流运行记录ID列表 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含删除结果的JSON响应 - """ - await WorkflowRunService.delete_service(auth=auth, ids=ids) - return SuccessResponse(data=None, msg="删除工作流运行记录成功") - - -@WorkflowRouter.delete( - "/run/clear", - summary="清空工作流运行记录", - description="清空工作流运行记录", - response_model=ResponseSchema[None], -) -async def clear_workflow_run_controller( - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:delete"]))], -) -> JSONResponse: - """ - 清空工作流运行记录 - - 参数: - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含清空结果的JSON响应 - """ - await WorkflowRunService.clear_service(auth=auth) - return SuccessResponse(data=None, msg="清空工作流运行记录成功") - - -@WorkflowRouter.post( - "/run/cancel/{id}", - summary="取消工作流运行记录", - description="取消工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def cancel_workflow_run_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], -) -> JSONResponse: - """ - 取消工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含取消的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.cancel_service(auth=auth, id=id) - return SuccessResponse(data=result_dict, msg="取消工作流运行记录成功") - - -@WorkflowRouter.post( - "/run/retry/{id}", - summary="重试工作流运行记录", - description="重试工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def retry_workflow_run_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], - retry_count: int = 1, -) -> JSONResponse: - """ - 重试工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - retry_count (int): 重试次数 - - 返回: - - JSONResponse: 包含重试的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.retry_service(auth=auth, id=id, retry_count=retry_count) - return SuccessResponse(data=result_dict, msg="重试工作流运行记录成功") - - -@WorkflowRouter.post( - "/run/pause/{id}", - summary="暂停工作流运行记录", - description="暂停工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def pause_workflow_run_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], -) -> JSONResponse: - """ - 暂停工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含暂停的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.pause_service(auth=auth, id=id) - return SuccessResponse(data=result_dict, msg="暂停工作流运行记录成功") - - -@WorkflowRouter.post( - "/run/resume/{id}", - summary="恢复工作流运行记录", - description="恢复工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def resume_workflow_run_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], -) -> JSONResponse: - """ - 恢复工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含恢复的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.resume_service(auth=auth, id=id) - return SuccessResponse(data=result_dict, msg="恢复工作流运行记录成功") - - -@WorkflowRouter.post( - "/run/terminate/{id}", - summary="终止工作流运行记录", - description="终止工作流运行记录", - response_model=ResponseSchema[WorkflowRunOutSchema], -) -async def terminate_workflow_run_controller( - id: int, - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:update"]))], -) -> JSONResponse: - """ - 终止工作流运行记录 - - 参数: - - id (int): 工作流运行记录ID - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含终止的工作流运行记录的JSON响应 - """ - result_dict = await WorkflowRunService.terminate_service(auth=auth, id=id) - return SuccessResponse(data=result_dict, msg="终止工作流运行记录成功") - - -@WorkflowRouter.get( - "/run/log/list", - summary="查询工作流运行日志列表", - description="查询工作流运行日志列表", - response_model=ResponseSchema[list[WorkflowRunLogOutSchema]], -) -async def get_workflow_run_log_list_controller( - page: Annotated[PaginationQueryParam, Depends()], - search: Annotated[WorkflowRunLogQueryParam, Depends()], - auth: Annotated[AuthSchema, Depends(AuthPermission(["module_task:workflow:run:query"]))], -) -> JSONResponse: - """ - 查询工作流运行日志列表 - - 参数: - - page (PaginationQueryParam): 分页查询参数 - - search (WorkflowRunLogQueryParam): 查询参数 - - auth (AuthSchema): 认证信息模型 - - 返回: - - JSONResponse: 包含工作流运行日志列表的JSON响应 - """ - result_dict = await WorkflowRunLogService.page_service( - auth=auth, - page_no=page.page_no, - page_size=page.page_size, - search=search.__dict__ if search else None, - order_by=page.order_by, - ) - return SuccessResponse(data=result_dict, msg="查询工作流运行日志列表成功") diff --git a/backend/app/plugin/module_task/workflow/crud.py b/backend/app/plugin/module_task/workflow/crud.py index 8d061bf9..b38e9921 100644 --- a/backend/app/plugin/module_task/workflow/crud.py +++ b/backend/app/plugin/module_task/workflow/crud.py @@ -1,18 +1,12 @@ 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 .model import WorkflowModel, WorkflowRunLogModel, WorkflowRunModel +from .model import WorkflowModel from .schema import ( WorkflowCreateSchema, WorkflowOutSchema, - WorkflowRunCreateSchema, - WorkflowRunLogCreateSchema, - WorkflowRunLogOutSchema, - WorkflowRunOutSchema, - WorkflowRunUpdateSchema, WorkflowUpdateSchema, ) @@ -156,124 +150,3 @@ class WorkflowCRUD(CRUDBase[WorkflowModel, WorkflowCreateSchema, WorkflowUpdateS out_schema=WorkflowOutSchema, preload=preload, ) - - async def list_templates_crud( - self, - search: dict | None = None, - order_by: list[dict] | None = None, - preload: list[str] | None = None, - ) -> Sequence[WorkflowModel]: - """ - 获取模板列表 - - 参数: - - search (dict | None): 查询参数 - - order_by (list[dict] | None): 排序参数 - - preload (list[str] | None): 预加载关系 - - 返回: - - Sequence[WorkflowModel]: 工作流模型实例序列 - """ - search_dict = search or {} - search_dict["is_template"] = (None, True) - return await self.list(search=search_dict, order_by=order_by, preload=preload) - - -class WorkflowRunCRUD(CRUDBase[WorkflowRunModel, WorkflowRunCreateSchema, WorkflowRunUpdateSchema]): - """工作流运行记录CRUD""" - - def __init__(self, auth: AuthSchema) -> None: - """初始化CRUD数据层""" - super().__init__(model=WorkflowRunModel, auth=auth) - - async def list_crud( - self, - search: dict | None = None, - order_by: list[dict] | None = None, - preload: list[str | Any] | None = None, - ) -> Sequence[WorkflowRunModel]: - """获取工作流运行记录列表""" - return await self.list(search=search, order_by=order_by, preload=preload) - - async def create_crud(self, data: WorkflowRunCreateSchema) -> WorkflowRunModel | None: - """创建工作流运行记录""" - return await self.create(data=data) - - async def get_by_id_crud(self, id: int) -> WorkflowRunModel | None: - """根据ID获取工作流运行记录""" - return await self.get(id=id) - - async def update_crud(self, id: int, data: WorkflowRunUpdateSchema) -> WorkflowRunModel | None: - """更新工作流运行记录""" - return await self.update(id=id, data=data) - - async def delete_crud(self, ids: list[int]) -> None: - """删除工作流运行记录""" - return await self.delete(ids=ids) - - async def clear_crud(self) -> None: - """清除所有工作流运行记录""" - return await self.clear() - - async def page_crud( - self, - offset: int, - limit: int, - order_by: list[dict] | None = None, - search: dict | None = None, - preload: list | None = None, - ) -> dict: - """分页查询工作流运行记录""" - order_by_list = order_by or [{"id": "desc"}] - search_dict = search or {} - - return await self.page( - offset=offset, - limit=limit, - order_by=order_by_list, - search=search_dict, - out_schema=WorkflowRunOutSchema, - preload=preload, - ) - - -class WorkflowRunLogCRUD(CRUDBase[WorkflowRunLogModel, WorkflowRunLogCreateSchema, WorkflowRunLogCreateSchema]): - """工作流运行日志CRUD""" - - def __init__(self, auth: AuthSchema) -> None: - """初始化CRUD数据层""" - super().__init__(model=WorkflowRunLogModel, auth=auth) - - async def list_crud( - self, - search: dict | None = None, - order_by: list[dict] | None = None, - preload: list[str | Any] | None = None, - ) -> Sequence[WorkflowRunLogModel]: - """获取工作流运行日志列表""" - return await self.list(search=search, order_by=order_by, preload=preload) - - async def create_crud(self, data: WorkflowRunLogCreateSchema) -> WorkflowRunLogModel | None: - """创建工作流运行日志""" - return await self.create(data=data) - - async def page_crud( - self, - offset: int, - limit: int, - order_by: list[dict] | None = None, - search: dict | None = None, - preload: list | None = None, - ) -> dict: - """分页查询工作流运行日志""" - order_by_list = order_by or [{"id": "desc"}] - search_dict = search or {} - - return await self.page( - offset=offset, - limit=limit, - order_by=order_by_list, - search=search_dict, - out_schema=WorkflowRunLogOutSchema, - preload=preload, - ) diff --git a/backend/app/plugin/module_task/workflow/engine.py b/backend/app/plugin/module_task/workflow/engine.py new file mode 100644 index 00000000..29fb663a --- /dev/null +++ b/backend/app/plugin/module_task/workflow/engine.py @@ -0,0 +1,397 @@ +import asyncio +import json +from collections import defaultdict +from collections.abc import Callable +from datetime import datetime +from typing import Any + +from sqlalchemy import select + +from app.core.database import async_db_session +from app.core.logger import log +from app.plugin.module_task.node.model import NodeModel +from app.plugin.module_task.workflow.model import WorkflowModel + + +class WorkflowRunLogModel: + pass + + +class NodeExecutionContext: + def __init__( + self, + node_id: str, + node_type: str, + args: str | None = None, + kwargs: str | None = None, + variables: dict | None = None, + input_data: Any = None, + ): + self.node_id = node_id + self.node_type = node_type + self.args = args + self.kwargs = kwargs + self.variables = variables or {} + self.input_data = input_data + self.output_data: Any = None + self.status: str = "pending" + self.error: str | None = None + self.start_time: datetime | None = None + self.end_time: datetime | None = None + + +class WorkflowExecutionContext: + def __init__( + self, + workflow_id: int, + workflow_name: str, + variables: dict | None = None, + business_key: str | None = None, + ): + self.workflow_id = workflow_id + self.workflow_name = workflow_name + self.variables = variables or {} + self.business_key = business_key + self.node_contexts: dict[str, NodeExecutionContext] = {} + self.status: str = "pending" + self.start_time: datetime | None = None + self.end_time: datetime | None = None + self.current_node: str | None = None + + +class WorkflowEngine: + """ + 工作流执行引擎 + + 功能: + 1. 解析工作流节点和边,构建执行图 + 2. 按拓扑顺序执行节点 + 3. 支持条件分支和并行执行 + 4. 记录执行日志 + """ + + _node_handlers: dict[str, Callable] = {} + + @classmethod + def register_handler(cls, node_type: str, handler: Callable) -> None: + """ + 注册节点处理器 + + 参数: + - node_type: 节点类型编码 + - handler: 处理函数,签名为 async (context: NodeExecutionContext) -> Any + """ + cls._node_handlers[node_type] = handler + + @classmethod + async def execute(cls, workflow: WorkflowModel, variables: dict | None = None, business_key: str | None = None) -> WorkflowExecutionContext: + """ + 执行工作流 + + 参数: + - workflow: 工作流模型实例 + - variables: 流程变量 + - business_key: 业务键 + + 返回: + - WorkflowExecutionContext: 执行上下文 + """ + context = WorkflowExecutionContext( + workflow_id=workflow.id, + workflow_name=workflow.name, + variables=variables or {}, + business_key=business_key, + ) + + context.start_time = datetime.now() + context.status = "running" + + try: + nodes = workflow.nodes if isinstance(workflow.nodes, list) else [] + edges = workflow.edges if isinstance(workflow.edges, list) else [] + + if not nodes: + context.status = "completed" + context.end_time = datetime.now() + return context + + node_map = {n.get("id"): n for n in nodes} + execution_order = cls._build_execution_order(nodes, edges) + + node_type_map = await cls._load_node_types([n.get("type") for n in nodes if n.get("type")]) + + for node_id in execution_order: + node_data = node_map.get(node_id) + if not node_data: + continue + + node_type_code = node_data.get("type") + node_args = node_data.get("data", {}).get("args", "") + node_kwargs = node_data.get("data", {}).get("kwargs", "{}") + + input_data = cls._collect_input_data(node_id, edges, context) + + node_context = NodeExecutionContext( + node_id=node_id, + node_type=node_type_code, + args=node_args, + kwargs=node_kwargs, + variables=context.variables, + input_data=input_data, + ) + context.node_contexts[node_id] = node_context + context.current_node = node_id + + node_type_info = node_type_map.get(node_type_code) + await cls._execute_node(node_context, node_type_info) + + if node_context.status == "failed": + context.status = "failed" + context.end_time = datetime.now() + return context + + if node_context.output_data: + context.variables[node_id] = node_context.output_data + + context.status = "completed" + except Exception as e: + log.error(f"工作流 {workflow.id} 执行失败: {e!s}") + context.status = "failed" + finally: + context.end_time = datetime.now() + context.current_node = None + + return context + + @classmethod + def _build_execution_order(cls, nodes: list, edges: list) -> list[str]: + """ + 构建节点执行顺序(拓扑排序) + + 参数: + - nodes: 节点列表 + - edges: 边列表 + + 返回: + - list[str]: 节点ID执行顺序 + """ + in_degree = defaultdict(int) + adjacency = defaultdict(list) + node_ids = {n.get("id") for n in nodes} + + for edge in edges: + source = edge.get("source") + target = edge.get("target") + if source in node_ids and target in node_ids: + adjacency[source].append(target) + in_degree[target] += 1 + + queue = [nid for nid in node_ids if in_degree[nid] == 0] + result = [] + + while queue: + node_id = queue.pop(0) + result.append(node_id) + + for neighbor in adjacency[node_id]: + in_degree[neighbor] -= 1 + if in_degree[neighbor] == 0: + queue.append(neighbor) + + for nid in node_ids: + if nid not in result: + result.append(nid) + + return result + + @classmethod + async def _load_node_types(cls, type_codes: list[str]) -> dict[str, NodeModel]: + """ + 加载节点类型定义 + + 参数: + - type_codes: 节点类型编码列表 + + 返回: + - dict[str, NodeModel]: 节点类型映射 + """ + if not type_codes: + return {} + + async with async_db_session() as session: + result = await session.execute(select(NodeModel).where(NodeModel.code.in_(type_codes))) + node_types = result.scalars().all() + return {nt.code: nt for nt in node_types} + + @classmethod + def _collect_input_data(cls, node_id: str, edges: list, context: WorkflowExecutionContext) -> dict: + """ + 收集节点输入数据(从上游节点的输出) + + 参数: + - node_id: 当前节点ID + - edges: 边列表 + - context: 执行上下文 + + 返回: + - dict: 输入数据 + """ + input_data = {} + for edge in edges: + if edge.get("target") == node_id: + source_id = edge.get("source") + source_context = context.node_contexts.get(source_id) + if source_context and source_context.output_data: + edge_label = edge.get("label", "default") + input_data[edge_label] = source_context.output_data + return input_data + + @classmethod + async def _execute_node(cls, context: NodeExecutionContext, node_type_info: NodeModel | None) -> None: + """ + 执行单个节点 + + 参数: + - context: 节点执行上下文 + - node_type_info: 节点类型定义 + """ + context.start_time = datetime.now() + context.status = "running" + + try: + handler = cls._node_handlers.get(context.node_type) + + if handler: + context.output_data = await handler(context) + elif node_type_info and node_type_info.func: + context.output_data = await cls._execute_code_block( + node_type_info.func, + node_type_info.args, + node_type_info.kwargs, + context, + ) + else: + log.warning(f"节点 {context.node_id} 没有注册处理器或代码块,跳过执行") + context.output_data = None + + context.status = "completed" + except Exception as e: + log.error(f"节点 {context.node_id} 执行失败: {e!s}") + context.status = "failed" + context.error = str(e) + finally: + context.end_time = datetime.now() + + @classmethod + async def _execute_code_block( + cls, + code_block: str, + args: str | None, + kwargs: str | None, + context: NodeExecutionContext, + ) -> Any: + """ + 执行代码块 + + 参数: + - code_block: 代码块 + - args: 位置参数 + - kwargs: 关键字参数 + - context: 节点执行上下文 + + 返回: + - Any: 执行结果 + """ + if not code_block or not code_block.strip(): + return None + + local_vars = { + "context": context, + "variables": context.variables, + "input_data": context.input_data, + } + + exec(code_block, {"__builtins__": __builtins__}, local_vars) + + handler = local_vars.get("handler") + if handler and callable(handler): + job_args = [] + if args: + args_str = str(args).strip() + if args_str: + job_args = [arg.strip() for arg in args_str.split(",") if arg.strip()] + + job_kwargs = {} + if kwargs: + kwargs_str = str(kwargs).strip() + if kwargs_str: + try: + job_kwargs = json.loads(kwargs_str) + except json.JSONDecodeError: + pass + + result = handler(*job_args, **job_kwargs) + if asyncio.iscoroutine(result): + return await result + return result + + return local_vars.get("result") + + +def register_builtin_handlers(): + """ + 注册内置节点处理器 + """ + + async def input_handler(context: NodeExecutionContext) -> Any: + return context.variables + + async def output_handler(context: NodeExecutionContext) -> Any: + return context.input_data + + async def condition_handler(context: NodeExecutionContext) -> Any: + try: + kwargs_str = str(context.kwargs or "{}").strip() + kwargs_data = json.loads(kwargs_str) if kwargs_str else {} + condition_expr = kwargs_data.get("condition", "True") + result = eval(condition_expr, {"__builtins__": __builtins__}, context.variables) + return {"condition_result": bool(result)} + except Exception as e: + log.error(f"条件表达式执行失败: {e!s}") + return {"condition_result": False} + + async def http_request_handler(context: NodeExecutionContext) -> Any: + import httpx + + kwargs_str = str(context.kwargs or "{}").strip() + kwargs_data = json.loads(kwargs_str) if kwargs_str else {} + + url = kwargs_data.get("url") + method = kwargs_data.get("method", "GET").upper() + headers = kwargs_data.get("headers", {}) + body = kwargs_data.get("body") + + if not url: + raise ValueError("HTTP请求节点缺少URL配置") + + async with httpx.AsyncClient() as client: + response = await client.request( + method=method, + url=url, + headers=headers, + json=body if body else None, + timeout=30.0, + ) + return { + "status_code": response.status_code, + "body": response.text, + "headers": dict(response.headers), + } + + WorkflowEngine.register_handler("input", input_handler) + WorkflowEngine.register_handler("output", output_handler) + WorkflowEngine.register_handler("condition", condition_handler) + WorkflowEngine.register_handler("http_request", http_request_handler) + + +register_builtin_handlers() diff --git a/backend/app/plugin/module_task/workflow/model.py b/backend/app/plugin/module_task/workflow/model.py index d938e7cf..3604be95 100644 --- a/backend/app/plugin/module_task/workflow/model.py +++ b/backend/app/plugin/module_task/workflow/model.py @@ -1,26 +1,11 @@ import enum -from datetime import datetime -from sqlalchemy import ( - JSON, - Boolean, - DateTime, - Integer, - String, - Text, -) +from sqlalchemy import JSON, String from sqlalchemy.orm import Mapped, mapped_column from app.core.base_model import ModelMixin, UserMixin -class StatusEnum(enum.Enum): - """状态枚举""" - - ACTIVE = "active" - INACTIVE = "inactive" - - class WorkflowStatusEnum(enum.Enum): """流程状态枚举""" @@ -29,18 +14,6 @@ class WorkflowStatusEnum(enum.Enum): ARCHIVED = "archived" -class RunStatusEnum(enum.Enum): - """运行状态枚举""" - - PENDING = "pending" - RUNNING = "running" - PAUSED = "paused" - COMPLETED = "completed" - FAILED = "failed" - CANCELLED = "cancelled" - TERMINATED = "terminated" - - class WorkflowModel(ModelMixin, UserMixin): """ 工作流模型 - 通用流程编排引擎 @@ -53,57 +26,5 @@ class WorkflowModel(ModelMixin, UserMixin): name: Mapped[str] = mapped_column(String(128), nullable=False, index=True, comment="流程名称") code: Mapped[str] = mapped_column(String(64), nullable=False, unique=True, index=True, comment="流程编码") status: Mapped[str] = mapped_column(String(32), default=WorkflowStatusEnum.DRAFT.value, nullable=False, index=True, comment="流程状态") - nodes: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict, comment="节点数据(JSON格式)") - edges: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict, comment="连线数据(JSON格式)") - version: Mapped[str] = mapped_column(String(32), default="1.0.0", comment="版本号") - category: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True, comment="流程分类") - tags: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="标签(JSON数组)") - is_template: Mapped[bool] = mapped_column(Boolean, default=False, index=True, comment="是否为模板") - template_id: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="模板ID") - published_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, index=True, comment="发布时间") - published_by: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="发布人ID") - meta_data: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="元数据(JSON格式)") - - canvas_state: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="画布状态(缩放、平移等)") - thumbnail: Mapped[str | None] = mapped_column(Text, nullable=True, comment="流程缩略图(Base64)") - statistics: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="流程统计信息") - - execution_config: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="执行配置(超时、重试等)") - variables: Mapped[dict | None] = mapped_column(JSON, nullable=True, comment="全局变量") - timeout: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="流程执行超时时间(秒)") - - -class WorkflowRunModel(ModelMixin, UserMixin): - """工作流运行记录模型""" - __tablename__: str = "task_workflow_run" - __table_args__: dict[str, str] = {"comment": "工作流运行记录表"} - - workflow_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True, comment="工作流ID") - workflow_name: Mapped[str] = mapped_column(String(128), nullable=False, comment="工作流名称") - workflow_version: Mapped[str] = mapped_column(String(32), default="1.0.0", comment="工作流版本") - business_key: Mapped[str | None] = mapped_column(String(128), nullable=True, comment="业务键") - initiator: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="发起人ID") - initiator_name: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="发起人姓名") - variables: Mapped[dict] = mapped_column(JSON, default={}, comment="流程变量(JSON格式)") - status: Mapped[str] = mapped_column(String(32), default=RunStatusEnum.PENDING.value, nullable=False, index=True, comment="运行状态") - error_message: Mapped[str | None] = mapped_column(Text, nullable=True, comment="错误信息") - start_time: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, comment="开始时间") - end_time: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, comment="结束时间") - duration: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="执行时长(秒)") - retry_count: Mapped[int] = mapped_column(Integer, default=0, comment="重试次数") - max_retry: Mapped[int] = mapped_column(Integer, default=3, comment="最大重试次数") - job_id: Mapped[int | None] = mapped_column(Integer, nullable=True, comment="关联的定时任务ID") - meta_data: Mapped[dict] = mapped_column(JSON, default={}, comment="元数据(JSON格式)") - - -class WorkflowRunLogModel(ModelMixin, UserMixin): - """工作流运行日志模型""" - __tablename__: str = "task_workflow_run_log" - __table_args__: dict[str, str] = {"comment": "工作流运行日志表"} - - task_run_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True, comment="工作流运行记录ID") - level: Mapped[str] = mapped_column(String(16), default="info", comment="日志级别") - node_id: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="节点ID") - node_name: Mapped[str | None] = mapped_column(String(128), nullable=True, comment="节点名称") - message: Mapped[str] = mapped_column(Text, nullable=False, comment="日志消息") - data: Mapped[dict] = mapped_column(JSON, default={}, comment="日志数据(JSON格式)") + nodes: Mapped[list] = mapped_column(JSON, nullable=False, default=list, comment="节点数据(JSON格式)") + edges: Mapped[list] = mapped_column(JSON, nullable=False, default=list, comment="连线数据(JSON格式)") diff --git a/backend/app/plugin/module_task/workflow/schema.py b/backend/app/plugin/module_task/workflow/schema.py index d81bd655..2c04eedd 100644 --- a/backend/app/plugin/module_task/workflow/schema.py +++ b/backend/app/plugin/module_task/workflow/schema.py @@ -1,12 +1,10 @@ from dataclasses import dataclass -from datetime import datetime from fastapi import Query from pydantic import ( BaseModel, ConfigDict, Field, - field_serializer, field_validator, model_validator, ) @@ -25,15 +23,6 @@ class WorkflowCreateSchema(BaseModel): description: str | None = Field(default=None, description="流程描述") nodes: list = Field(default_factory=list, description="节点数据(JSON格式)") edges: list = Field(default_factory=list, description="连线数据(JSON格式)") - version: str = Field(default="1.0.0", description="版本号") - category: str | None = Field(default=None, description="流程分类") - tags: list[str] | None = Field(default=None, description="标签列表") - is_template: bool = Field(default=False, description="是否为模板") - template_id: int | None = Field(default=None, description="模板ID") - meta_data: dict | None = Field(default=None, description="元数据(JSON格式)") - canvas_state: dict | None = Field(default=None, description="画布状态(缩放、平移等)") - thumbnail: str | None = Field(default=None, description="流程缩略图(Base64)") - statistics: dict | None = Field(default=None, description="流程统计信息") @field_validator("code") @classmethod @@ -58,8 +47,6 @@ class WorkflowCreateSchema(BaseModel): @model_validator(mode="after") def _after_validation(self): """核心业务规则校验""" - if self.is_template and self.template_id: - raise ValueError("模板不能引用其他模板") if self.description and len(self.description) > 1000: raise ValueError("流程描述长度不能超过1000个字符") return self @@ -74,17 +61,6 @@ class WorkflowUpdateSchema(BaseModel): description: str | None = Field(default=None, description="流程描述") nodes: list | None = Field(default=None, description="节点数据(JSON格式)") edges: list | None = Field(default=None, description="连线数据(JSON格式)") - version: str | None = Field(default=None, description="版本号") - category: str | None = Field(default=None, description="流程分类") - tags: list[str] | None = Field(default=None, description="标签列表") - is_template: bool | None = Field(default=None, description="是否为模板") - template_id: int | None = Field(default=None, description="模板ID") - published_at: str | None = Field(default=None, description="发布时间") - published_by: int | None = Field(default=None, description="发布人ID") - meta_data: dict | None = Field(default=None, description="元数据(JSON格式)") - canvas_state: dict | None = Field(default=None, description="画布状态(缩放、平移等)") - thumbnail: str | None = Field(default=None, description="流程缩略图(Base64)") - statistics: dict | None = Field(default=None, description="流程统计信息") @field_validator("status") @classmethod @@ -99,8 +75,6 @@ class WorkflowUpdateSchema(BaseModel): @model_validator(mode="after") def _after_validation(self): """核心业务规则校验""" - if self.is_template and self.template_id: - raise ValueError("模板不能引用其他模板") if self.description and len(self.description) > 1000: raise ValueError("流程描述长度不能超过1000个字符") return self @@ -115,23 +89,7 @@ class WorkflowOutSchema(WorkflowUpdateSchema, BaseSchema, UserBySchema): class WorkflowPublishSchema(BaseModel): """发布工作流模型""" - version: str | None = Field(default=None, description="版本号") - - -class WorkflowValidateSchema(BaseModel): - """验证工作流模型""" - - nodes: list = Field(..., description="节点数据(JSON格式)") - edges: list = Field(..., description="连线数据(JSON格式)") - - -class WorkflowValidateResultSchema(BaseModel): - """工作流验证结果模型""" - - is_valid: bool = Field(..., description="是否有效") - errors: list[str] = Field(default_factory=list, description="错误列表") - warnings: list[str] = Field(default_factory=list, description="警告列表") - stats: dict = Field(default_factory=dict, description="统计信息") + pass @dataclass @@ -143,8 +101,6 @@ class WorkflowQueryParam: name: str | None = Query(None, description="流程名称"), code: str | None = Query(None, description="流程编码"), status: str | None = Query(None, description="流程状态"), - category: str | None = Query(None, description="流程分类"), - is_template: bool | None = Query(None, description="是否为模板"), created_time: list[DateTimeStr] | None = Query( None, description="创建时间范围", @@ -164,10 +120,6 @@ class WorkflowQueryParam: self.code = (QueueEnum.like.value, code) if status: self.status = (QueueEnum.eq.value, status) - if category: - self.category = (QueueEnum.eq.value, category) - if is_template is not None: - self.is_template = (QueueEnum.eq.value, is_template) 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: @@ -190,124 +142,10 @@ class WorkflowExecuteSchema(BaseModel): class WorkflowExecuteResultSchema(BaseModel): """工作流执行结果模型""" - task_run_id: int = Field(..., description="任务运行记录ID") workflow_id: int = Field(..., description="工作流ID") workflow_name: str = Field(..., description="工作流名称") status: str = Field(..., description="执行状态") - message: str = Field(..., description="执行消息") - - -class WorkflowRunCreateSchema(BaseModel): - """创建工作流运行记录模型""" - workflow_id: int = Field(..., description="工作流ID") - workflow_name: str = Field(..., description="工作流名称") - workflow_version: str = Field(default="1.0.0", description="工作流版本") - business_key: str | None = Field(default=None, description="业务键") - initiator: int | None = Field(default=None, description="发起人ID") - initiator_name: str | None = Field(default=None, description="发起人姓名") - variables: dict = Field(default_factory=dict, description="流程变量(JSON格式)") - job_id: int | None = Field(default=None, description="关联的定时任务ID") - - -class WorkflowRunUpdateSchema(BaseModel): - """更新工作流运行记录模型""" - status: str | None = Field(default=None, description="运行状态") - error_message: str | None = Field(default=None, description="错误信息") - start_time: datetime | None = Field(default=None, description="开始时间") - end_time: datetime | None = Field(default=None, description="结束时间") - duration: int | None = Field(default=None, description="执行时长(秒)") - retry_count: int | None = Field(default=None, description="重试次数") - max_retry: int | None = Field(default=None, description="最大重试次数") - meta_data: dict | None = Field(default=None, description="元数据(JSON格式)") - - -class WorkflowRunOutSchema(WorkflowRunCreateSchema, WorkflowRunUpdateSchema, BaseSchema, UserBySchema): - """工作流运行记录响应模型""" - - model_config = ConfigDict(from_attributes=True) - - @field_serializer('start_time', 'end_time') - def serialize_datetime(self, value: datetime | None) -> str | None: - if value is None: - return None - return value.strftime("%Y-%m-%d %H:%M:%S") - - -class WorkflowRunLogCreateSchema(BaseModel): - """创建工作流运行日志模型""" - task_run_id: int = Field(..., description="工作流运行记录ID") - level: str = Field(default="info", description="日志级别") - node_id: str | None = Field(default=None, description="节点ID") - node_name: str | None = Field(default=None, description="节点名称") - message: str = Field(..., description="日志消息") - data: dict | None = Field(default=None, description="日志数据(JSON格式)") - - -class WorkflowRunLogOutSchema(BaseModel): - """工作流运行日志响应模型""" - id: int - task_run_id: int - level: str - node_id: str | None - node_name: str | None - message: str - data: dict | None - created_time: str - created_by: dict | None - - model_config = ConfigDict(from_attributes=True) - - -@dataclass -class WorkflowRunQueryParam: - """工作流运行记录查询参数""" - def __init__( - self, - workflow_id: int | None = Query(None, description="工作流ID"), - workflow_name: str | None = Query(None, description="工作流名称"), - status: str | None = Query(None, description="运行状态"), - business_key: str | None = Query(None, description="业务键"), - initiator: int | None = Query(None, description="发起人"), - job_id: int | None = Query(None, description="定时任务ID"), - created_time: list[str] | None = Query(None, description="创建时间范围"), - updated_time: list[str] | None = Query(None, description="更新时间范围"), - ) -> None: - if workflow_id: - self.workflow_id = (QueueEnum.eq.value, workflow_id) - if workflow_name: - self.workflow_name = (QueueEnum.like.value, workflow_name) - if status: - self.status = (QueueEnum.eq.value, status) - if business_key: - self.business_key = (QueueEnum.like.value, business_key) - if initiator: - self.initiator = (QueueEnum.eq.value, initiator) - if job_id: - self.job_id = (QueueEnum.eq.value, job_id) - 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])) - - -@dataclass -class WorkflowRunLogQueryParam: - """工作流运行日志查询参数""" - def __init__( - self, - task_run_id: int | None = Query(None, description="工作流运行记录ID"), - level: str | None = Query(None, description="日志级别"), - node_id: str | None = Query(None, description="节点ID"), - node_name: str | None = Query(None, description="节点名称"), - created_time: list[str] | None = Query(None, description="创建时间范围"), - ) -> None: - if task_run_id: - self.task_run_id = (QueueEnum.eq.value, task_run_id) - if level: - self.level = (QueueEnum.eq.value, level) - if node_id: - self.node_id = (QueueEnum.like.value, node_id) - if node_name: - self.node_name = (QueueEnum.like.value, node_name) - if created_time and len(created_time) == 2: - self.created_time = (QueueEnum.between.value, (created_time[0], created_time[1])) + start_time: str | None = Field(default=None, description="开始时间") + end_time: str | None = Field(default=None, description="结束时间") + variables: dict = Field(default_factory=dict, description="流程变量") + node_results: dict = Field(default_factory=dict, description="节点执行结果") diff --git a/backend/app/plugin/module_task/workflow/service.py b/backend/app/plugin/module_task/workflow/service.py index 8797bd3b..f2c706bb 100644 --- a/backend/app/plugin/module_task/workflow/service.py +++ b/backend/app/plugin/module_task/workflow/service.py @@ -1,24 +1,16 @@ -from datetime import datetime -from typing import Any from app.api.v1.module_system.auth.schema import AuthSchema from app.core.exceptions import CustomException -from .crud import WorkflowCRUD, WorkflowRunCRUD, WorkflowRunLogCRUD +from .crud import WorkflowCRUD +from .engine import WorkflowEngine from .schema import ( WorkflowCreateSchema, WorkflowExecuteSchema, WorkflowOutSchema, WorkflowPublishSchema, WorkflowQueryParam, - WorkflowRunCreateSchema, - WorkflowRunLogCreateSchema, - WorkflowRunLogOutSchema, - WorkflowRunOutSchema, - WorkflowRunUpdateSchema, WorkflowUpdateSchema, - WorkflowValidateResultSchema, - WorkflowValidateSchema, ) @@ -44,28 +36,6 @@ class WorkflowService: raise CustomException(msg="该工作流不存在") return WorkflowOutSchema.model_validate(obj).model_dump() - @classmethod - async def list_service( - cls, - auth: AuthSchema, - search: WorkflowQueryParam | None = None, - order_by: list[dict[str, str]] | None = None, - ) -> list[dict]: - """ - 列表查询 - - 参数: - - auth (AuthSchema): 认证信息模型 - - search (WorkflowQueryParam | None): 查询参数 - - order_by (list[dict[str, str]] | None): 排序参数 - - 返回: - - list[dict]: 工作流模型实例字典列表 - """ - search_dict = search.__dict__ if search else None - obj_list = await WorkflowCRUD(auth).list_crud(search=search_dict, order_by=order_by) - return [WorkflowOutSchema.model_validate(obj).model_dump() for obj in obj_list] - @classmethod async def page_service( cls, @@ -116,11 +86,6 @@ class WorkflowService: if obj: raise CustomException(msg="创建失败,流程编码已存在") - if data.template_id: - template = await WorkflowCRUD(auth).get_by_id_crud(id=data.template_id) - if not template or not template.is_template: - raise CustomException(msg="创建失败,模板不存在或不是模板") - obj = await WorkflowCRUD(auth).create_crud(data=data) return WorkflowOutSchema.model_validate(obj).model_dump() @@ -141,19 +106,11 @@ class WorkflowService: if not obj: raise CustomException(msg="更新失败,该工作流不存在") - if obj.status == "published": - raise CustomException(msg="更新失败,已发布的工作流不能直接修改,请先归档") - if data.code: exist_obj = await WorkflowCRUD(auth).get_by_code_crud(code=data.code) if exist_obj and exist_obj.id != id: raise CustomException(msg="更新失败,流程编码重复") - if data.template_id: - template = await WorkflowCRUD(auth).get_by_id_crud(id=data.template_id) - if not template or not template.is_template: - raise CustomException(msg="更新失败,模板不存在或不是模板") - obj = await WorkflowCRUD(auth).update_crud(id=id, data=data) return WorkflowOutSchema.model_validate(obj).model_dump() @@ -174,10 +131,6 @@ class WorkflowService: obj_list = await WorkflowCRUD(auth).list_crud(search={"id": (None, ids)}) - for obj in obj_list: - if obj.status == "published": - raise CustomException(msg=f"删除失败,ID为{obj.id}的工作流已发布,无法删除") - if len(obj_list) != len(ids): existing_ids = {obj.id for obj in obj_list} missing_ids = set(ids) - existing_ids @@ -207,137 +160,11 @@ class WorkflowService: update_data = WorkflowUpdateSchema( status="published", - version=data.version or obj.version, - published_at=datetime.now().isoformat(), - published_by=auth.user.id if auth.user else None, ) obj = await WorkflowCRUD(auth).update_crud(id=id, data=update_data) return WorkflowOutSchema.model_validate(obj).model_dump() - @classmethod - async def validate_service(cls, auth: AuthSchema, data: WorkflowValidateSchema) -> WorkflowValidateResultSchema: - """ - 验证工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - data (WorkflowValidateSchema): 验证模型 - - 返回: - - WorkflowValidateResultSchema: 验证结果 - """ - errors = [] - warnings = [] - stats = {} - - nodes = data.nodes if isinstance(data.nodes, list) else [] - edges = data.edges if isinstance(data.edges, list) else [] - - stats["totalNodes"] = len(nodes) - stats["totalEdges"] = len(edges) - stats["nodeTypes"] = {} - - for node in nodes: - node_type = node.get("type", "unknown") - stats["nodeTypes"][node_type] = stats["nodeTypes"].get(node_type, 0) + 1 - - if len(nodes) == 0: - errors.append("流程中没有节点") - - node_ids = {n.get("id") for n in nodes} - for edge in edges: - source = edge.get("source") - target = edge.get("target") - if source not in node_ids: - errors.append(f"连线 {edge.get('label') or edge.get('id')} 的源节点不存在") - if target not in node_ids: - errors.append(f"连线 {edge.get('label') or edge.get('id')} 的目标节点不存在") - - is_valid = len(errors) == 0 - - return WorkflowValidateResultSchema( - is_valid=is_valid, - errors=errors, - warnings=warnings, - stats=stats, - ) - - @classmethod - async def templates_service(cls, auth: AuthSchema) -> list[dict]: - """ - 获取模板列表 - - 参数: - - auth (AuthSchema): 认证信息模型 - - 返回: - - list[dict]: 模板列表 - """ - obj_list = await WorkflowCRUD(auth).list_templates_crud(order_by=[{"id": "asc"}]) - return [WorkflowOutSchema.model_validate(obj).model_dump() for obj in obj_list] - - @classmethod - async def export_service(cls, auth: AuthSchema, id: int) -> dict: - """ - 导出工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - id (int): 工作流ID - - 返回: - - dict: 工作流数据 - """ - obj = await WorkflowCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="导出失败,该工作流不存在") - - return { - "id": obj.id, - "name": obj.name, - "code": obj.code, - "description": obj.description, - "nodes": obj.nodes, - "edges": obj.edges, - "version": obj.version, - "category": obj.category, - "metadata": obj.meta_data, - "exportedAt": datetime.now().isoformat(), - } - - @classmethod - async def import_service(cls, auth: AuthSchema, data: dict) -> dict: - """ - 导入工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - data (dict): 工作流数据 - - 返回: - - dict: 工作流模型实例字典 - """ - code = data.get("code") - if code: - exist_obj = await WorkflowCRUD(auth).get_by_code_crud(code=code) - if exist_obj: - raise CustomException(msg="导入失败,流程编码已存在") - - create_data = WorkflowCreateSchema( - name=data.get("name", "导入的工作流"), - code=data.get("code", f"imported_{datetime.now().timestamp()}"), - description=data.get("description"), - nodes=data.get("nodes", []), - edges=data.get("edges", []), - version=data.get("version", "1.0.0"), - category=data.get("category"), - meta_data=data.get("metadata"), - ) - - obj = await WorkflowCRUD(auth).create_crud(data=create_data) - return WorkflowOutSchema.model_validate(obj).model_dump() - @classmethod async def execute_service(cls, auth: AuthSchema, data: WorkflowExecuteSchema) -> dict: """ @@ -350,713 +177,32 @@ class WorkflowService: 返回: - dict: 执行结果 """ - task_run_id = await WorkflowRunService.execute_workflow( - auth=auth, - workflow_id=data.workflow_id, - variables=data.variables, - business_key=data.business_key, - initiator=auth.user.id if auth.user else None, - initiator_name=auth.user.username if auth.user else None, - job_id=data.job_id, - ) - - return { - "task_run_id": task_run_id, - "workflow_id": data.workflow_id, - "status": "running", - "message": "任务执行已启动", - } - - @classmethod - async def pause_service(cls, auth: AuthSchema, id: int) -> None: - """ - 暂停工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - id (int): 工作流ID - - 返回: - - None - """ - obj = await WorkflowCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="工作流不存在") - - await WorkflowRunService.pause_service(auth=auth, id=id) - - @classmethod - async def resume_service(cls, auth: AuthSchema, id: int) -> None: - """ - 恢复工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - id (int): 工作流ID - - 返回: - - None - """ - obj = await WorkflowCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="工作流不存在") - - await WorkflowRunService.resume_service(auth=auth, id=id) - - @classmethod - async def terminate_service(cls, auth: AuthSchema, id: int) -> None: - """ - 终止工作流 - - 参数: - - auth (AuthSchema): 认证信息模型 - - id (int): 工作流ID - - 返回: - - None - """ - obj = await WorkflowCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="工作流不存在") - - await WorkflowRunService.terminate_service(auth=auth, id=id) - - -class WorkflowRunService: - """工作流运行记录管理模块服务层""" - - @classmethod - async def list_service(cls, auth: AuthSchema, search: dict | None = None, order_by: list[dict] | None = None) -> list[dict]: - """获取工作流运行记录列表""" - search_dict = search.__dict__ if search else {} - order_by_list = order_by or [{"id": "desc"}] - - obj_list = await WorkflowRunCRUD(auth).list_crud(search=search_dict, order_by=order_by_list) - return [WorkflowRunOutSchema.model_validate(obj).model_dump() for obj in obj_list] - - @classmethod - async def page_service( - cls, - auth: AuthSchema, - page_no: int, - page_size: int, - search: dict | None = None, - order_by: list[dict] | None = None, - ) -> dict: - """分页查询工作流运行记录""" - search_dict = search.__dict__ if search else {} - order_by_list = order_by or [{"id": "desc"}] - offset = (page_no - 1) * page_size - - result = await WorkflowRunCRUD(auth).page_crud( - offset=offset, - limit=page_size, - order_by=order_by_list, - search=search_dict, - ) - return result - - @classmethod - async def detail_service(cls, auth: AuthSchema, id: int) -> dict: - """获取工作流运行记录详情""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="查询失败,该工作流运行记录不存在") - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def create_service(cls, auth: AuthSchema, data: WorkflowRunCreateSchema) -> dict: - """创建工作流运行记录""" - obj = await WorkflowRunCRUD(auth).create_crud(data=data) - if not obj: - raise CustomException(msg="创建失败,工作流运行记录创建失败") - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def update_service(cls, auth: AuthSchema, id: int, data: WorkflowRunUpdateSchema) -> dict: - """更新工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="更新失败,该工作流运行记录不存在") - - obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=data) - if not obj: - raise CustomException(msg="更新失败,工作流运行记录更新失败") - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def delete_service(cls, auth: AuthSchema, ids: list[int]) -> None: - """删除工作流运行记录""" - if len(ids) < 1: - raise CustomException(msg="删除失败,删除对象不能为空") - - obj_list = await WorkflowRunCRUD(auth).list_crud(search={"id": (None, ids)}) - if not obj_list or len(obj_list) != len(ids): - raise CustomException(msg="删除失败,部分工作流运行记录不存在") - - await WorkflowRunCRUD(auth).delete_crud(ids=ids) - - @classmethod - async def clear_service(cls, auth: AuthSchema) -> None: - """清除所有工作流运行记录""" - await WorkflowRunCRUD(auth).clear_crud() - - @classmethod - async def cancel_service(cls, auth: AuthSchema, id: int) -> dict: - """取消工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="取消失败,该工作流运行记录不存在") - - if obj.status in ["completed", "failed", "cancelled"]: - raise CustomException(msg="取消失败,该工作流运行记录已完成或已取消") - - update_data = WorkflowRunUpdateSchema( - status="cancelled", - end_time=datetime.now(), - ) - obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=update_data) - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def pause_service(cls, auth: AuthSchema, id: int) -> dict: - """暂停工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="暂停失败,该工作流运行记录不存在") - - if obj.status not in ["running", "pending"]: - raise CustomException(msg="暂停失败,只能暂停运行中或待执行的工作流运行记录") - - update_data = WorkflowRunUpdateSchema( - status="paused", - ) - obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=update_data) - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def resume_service(cls, auth: AuthSchema, id: int) -> dict: - """恢复工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="恢复失败,该工作流运行记录不存在") - - if obj.status != "paused": - raise CustomException(msg="恢复失败,只能恢复已暂停的工作流运行记录") - - update_data = WorkflowRunUpdateSchema( - status="running", - ) - obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=update_data) - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def terminate_service(cls, auth: AuthSchema, id: int) -> dict: - """终止工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="终止失败,该工作流运行记录不存在") - - if obj.status in ["completed", "failed", "cancelled", "terminated"]: - raise CustomException(msg="终止失败,该工作流运行记录已完成或已终止") - - update_data = WorkflowRunUpdateSchema( - status="terminated", - end_time=datetime.now(), - ) - obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=update_data) - return WorkflowRunOutSchema.model_validate(obj).model_dump() - - @classmethod - async def retry_service(cls, auth: AuthSchema, id: int, retry_count: int = 1) -> dict: - """重试工作流运行记录""" - obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not obj: - raise CustomException(msg="重试失败,该工作流运行记录不存在") - - if obj.status not in ["failed", "cancelled"]: - raise CustomException(msg="重试失败,只能重试失败或已取消的工作流运行记录") - - if obj.retry_count >= obj.max_retry: - raise CustomException(msg=f"重试失败,已达到最大重试次数{obj.max_retry}") - - update_data = WorkflowRunUpdateSchema( - status="pending", - retry_count=obj.retry_count + retry_count, - error_message=None, - start_time=None, - end_time=None, - duration=None, - ) - updated_obj = await WorkflowRunCRUD(auth).update_crud(id=id, data=update_data) - - if not updated_obj: - raise CustomException(msg="重试失败,更新工作流运行记录失败") - - try: - await cls.execute_workflow( - auth=auth, - workflow_id=obj.workflow_id, - variables=obj.variables, - business_key=obj.business_key, - initiator=obj.initiator, - initiator_name=obj.initiator_name, - job_id=obj.job_id, - task_run_id=id, - ) - except Exception: - pass - - final_obj = await WorkflowRunCRUD(auth).get_by_id_crud(id=id) - if not final_obj: - raise CustomException(msg="重试失败,获取工作流运行记录失败") - - return WorkflowRunOutSchema.model_validate(final_obj).model_dump() - - @classmethod - async def execute_workflow( - cls, - auth: AuthSchema, - workflow_id: int, - variables: dict, - business_key: str | None, - initiator: int | None, - initiator_name: str | None, - job_id: int | None, - task_run_id: int | None = None, - ) -> int: - """执行工作流""" - workflow = await WorkflowCRUD(auth).get_by_id_crud(id=workflow_id) + workflow = await WorkflowCRUD(auth).get_by_id_crud(id=data.workflow_id) if not workflow: raise CustomException(msg="工作流不存在") - nodes = workflow.nodes if isinstance(workflow.nodes, list) else [] - edges = workflow.edges if isinstance(workflow.edges, list) else [] + if workflow.status != "published": + raise CustomException(msg="工作流未发布,无法执行") - if not nodes: - raise CustomException(msg="工作流没有节点") - - start_time = datetime.now() - - if not task_run_id: - create_data = WorkflowRunCreateSchema( - workflow_id=workflow_id, - workflow_name=workflow.name, - workflow_version=workflow.version, - business_key=business_key, - initiator=initiator, - initiator_name=initiator_name, - variables=variables, - job_id=job_id, - ) - task_run = await WorkflowRunCRUD(auth).create_crud(data=create_data) - if not task_run: - raise CustomException(msg="创建工作流运行记录失败") - task_run_id = task_run.id - - if task_run_id is None: - raise CustomException(msg="工作流运行记录ID不能为空") - - update_data = WorkflowRunUpdateSchema( - status="running", - start_time=start_time, - ) - await WorkflowRunCRUD(auth).update_crud(id=task_run_id, data=update_data) - - try: - await cls._execute_nodes( - auth=auth, - task_run_id=task_run_id, - nodes=nodes, - edges=edges, - variables=variables, - ) - - end_time = datetime.now() - duration = int((end_time - start_time).total_seconds()) - - update_data = WorkflowRunUpdateSchema( - status="completed", - end_time=end_time, - duration=duration, - ) - await WorkflowRunCRUD(auth).update_crud(id=task_run_id, data=update_data) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - message=f"工作流执行完成,总耗时: {duration}秒", - ) - - except Exception as e: - end_time = datetime.now() - duration = int((end_time - start_time).total_seconds()) - - update_data = WorkflowRunUpdateSchema( - status="failed", - end_time=end_time, - duration=duration, - error_message=str(e), - ) - await WorkflowRunCRUD(auth).update_crud(id=task_run_id, data=update_data) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="error", - message=f"工作流执行失败: {e!s}", - ) - - raise - - return task_run_id - - @classmethod - async def _execute_nodes( - cls, - auth: AuthSchema, - task_run_id: int, - nodes: list[dict[str, Any]], - edges: list[dict[str, Any]], - variables: dict, - ) -> None: - """执行节点""" - node_map = {node.get("id"): node for node in nodes} - edge_map = {} - target_nodes = set() - for edge in edges: - source = edge.get("source") - target = edge.get("target") - if source not in edge_map: - edge_map[source] = [] - edge_map[source].append(edge) - target_nodes.add(target) - - start_nodes = [node for node in nodes if node.get("id") not in target_nodes] - if not start_nodes: - start_nodes = [node for node in nodes if node.get("type") == "trigger"] - if not start_nodes: - raise CustomException(msg="工作流缺少起始节点") - - executed_nodes = set() - queue = start_nodes.copy() - - while queue: - current_node = queue.pop(0) - node_id = current_node.get("id") - - if node_id in executed_nodes: - continue - - executed_nodes.add(node_id) - - await cls._execute_node( - auth=auth, - task_run_id=task_run_id, - node=current_node, - variables=variables, - ) - - if node_id in edge_map: - for edge in edge_map[node_id]: - target_id = edge.get("target") - target_node = node_map.get(target_id) - if target_node and target_id not in executed_nodes: - target_node_type = target_node.get("type") - - if target_node_type == "condition": - condition = target_node.get("data", {}).get("condition", "True") - try: - result = bool(eval(condition, {}, variables)) - except Exception: - result = False - - for next_edge in edge_map.get(target_id, []): - next_target_id = next_edge.get("target") - next_target_node = node_map.get(next_target_id) - edge_label = next_edge.get("label", "") - - if next_target_node and next_target_id not in executed_nodes: - if edge_label == "true" and result: - queue.append(next_target_node) - elif edge_label == "false" and not result: - queue.append(next_target_node) - else: - queue.append(target_node) - - @classmethod - async def _execute_node( - cls, - auth: AuthSchema, - task_run_id: int, - node: dict[str, Any], - variables: dict, - ) -> None: - """执行单个节点""" - node_id = str(node.get("id", "")) - node_type = node.get("type", "") - node_name = node.get("data", {}).get("label", "未知节点") - node_config = node.get("data", {}) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"开始执行节点: {node_name}", + context = await WorkflowEngine.execute( + workflow=workflow, + variables=data.variables, + business_key=data.business_key, ) - try: - if node_type == "task": - await cls._execute_task_node(auth, task_run_id, node_id, node_name, node_config, variables) - elif node_type == "approval": - await cls._execute_approval_node(auth, task_run_id, node_id, node_name, node_config, variables) - elif node_type == "condition": - await cls._execute_condition_node(auth, task_run_id, node_id, node_name, node_config, variables) - elif node_type == "notification": - await cls._execute_notification_node(auth, task_run_id, node_id, node_name, node_config, variables) - elif node_type == "timer": - await cls._execute_timer_node(auth, task_run_id, node_id, node_name, node_config, variables) - elif node_type == "parallel": - await cls._execute_parallel_node(auth, task_run_id, node_id, node_name, node_config, variables) - else: - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="warning", - node_id=node_id, - node_name=node_name, - message=f"未知节点类型: {node_type}", - ) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"节点执行完成: {node_name}", - ) - - except Exception as e: - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="error", - node_id=node_id, - node_name=node_name, - message=f"节点执行失败: {node_name}, 错误: {e!s}", - ) - raise - - @classmethod - async def _execute_task_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行任务节点""" - task_type = node_config.get("task_type", "default") - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"执行任务类型: {task_type}", - ) - - variables[f"{node_id}_result"] = {"status": "success", "output": "任务执行成功"} - - @classmethod - async def _execute_approval_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行审批节点""" - approvers = node_config.get("approvers", []) - approval_type = node_config.get("approval_type", "or") - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"审批类型: {approval_type}, 审批人: {approvers}", - ) - - variables[f"{node_id}_result"] = {"status": "approved", "output": "审批通过"} - - @classmethod - async def _execute_condition_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行条件节点""" - condition = node_config.get("condition", "True") - - try: - result = bool(eval(condition, {}, variables)) - except Exception as e: - result = False - raise CustomException(msg=f"条件表达式执行失败: {e!s}") - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"条件表达式: {condition}, 结果: {result}", - ) - - variables[f"{node_id}_result"] = {"status": "success", "result": result} - - @classmethod - async def _execute_notification_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行通知节点""" - notification_type = node_config.get("notification_type", "email") - recipients = node_config.get("recipients", []) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"发送通知: {notification_type}, 收件人: {recipients}", - ) - - variables[f"{node_id}_result"] = {"status": "sent", "output": "通知发送成功"} - - @classmethod - async def _execute_timer_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行定时节点""" - timer_type = node_config.get("timer_type", "delay") - timer_value = node_config.get("timer_value", 0) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"定时器类型: {timer_type}, 值: {timer_value}", - ) - - variables[f"{node_id}_result"] = {"status": "completed", "output": "定时器执行完成"} - - @classmethod - async def _execute_parallel_node( - cls, - auth: AuthSchema, - task_run_id: int, - node_id: str, - node_name: str, - node_config: dict, - variables: dict, - ) -> None: - """执行并行节点""" - parallel_count = node_config.get("parallel_count", 1) - - await cls._add_log( - auth=auth, - task_run_id=task_run_id, - level="info", - node_id=node_id, - node_name=node_name, - message=f"并行执行数量: {parallel_count}", - ) - - variables[f"{node_id}_result"] = {"status": "completed", "output": "并行执行完成"} - - @classmethod - async def _add_log( - cls, - auth: AuthSchema, - task_run_id: int, - level: str, - message: str, - node_id: str | None = None, - node_name: str | None = None, - data: dict | None = None, - ) -> None: - """添加日志""" - log_data = WorkflowRunLogCreateSchema( - task_run_id=task_run_id, - level=level, - node_id=node_id, - node_name=node_name, - message=message, - data=data or {}, - ) - await WorkflowRunLogCRUD(auth).create_crud(data=log_data) - - -class WorkflowRunLogService: - """工作流运行日志管理模块服务层""" - - @classmethod - async def list_service(cls, auth: AuthSchema, search: dict | None = None, order_by: list[dict] | None = None) -> list[dict]: - """获取工作流运行日志列表""" - search_dict = search.__dict__ if search else {} - order_by_list = order_by or [{"id": "desc"}] - - obj_list = await WorkflowRunLogCRUD(auth).list_crud(search=search_dict, order_by=order_by_list) - return [WorkflowRunLogOutSchema.model_validate(obj).model_dump() for obj in obj_list] - - @classmethod - async def page_service( - cls, - auth: AuthSchema, - page_no: int, - page_size: int, - search: dict | None = None, - order_by: list[dict] | None = None, - ) -> dict: - """分页查询工作流运行日志""" - search_dict = search.__dict__ if search else {} - order_by_list = order_by or [{"id": "desc"}] - offset = (page_no - 1) * page_size - - result = await WorkflowRunLogCRUD(auth).page_crud( - offset=offset, - limit=page_size, - order_by=order_by_list, - search=search_dict, - ) - return result + return { + "workflow_id": data.workflow_id, + "workflow_name": workflow.name, + "status": context.status, + "start_time": context.start_time.isoformat() if context.start_time else None, + "end_time": context.end_time.isoformat() if context.end_time else None, + "variables": context.variables, + "node_results": { + node_id: { + "status": ctx.status, + "output": ctx.output_data, + "error": ctx.error, + } + for node_id, ctx in context.node_contexts.items() + }, + } diff --git a/backend/app/scripts/data/sys_menu.json b/backend/app/scripts/data/sys_menu.json index 5e9f621f..08f16d6a 100644 --- a/backend/app/scripts/data/sys_menu.json +++ b/backend/app/scripts/data/sys_menu.json @@ -154,7 +154,7 @@ "description": "初始化数据" }, { - "name": "详情改菜", + "name": "详情菜单", "type": 3, "icon": null, "order": 5, @@ -166,7 +166,7 @@ "keep_alive": true, "hidden": false, "always_show": false, - "title": "详情改菜", + "title": "详情菜单", "params": null, "affix": false, "redirect": null, @@ -2901,68 +2901,11 @@ "redirect": null, "description": "发布工作流" }, - { - "name": "验证工作流", - "type": 3, - "icon": null, - "order": 5, - "permission": "module_task:workflow:validate", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "验证工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "验证工作流" - }, - { - "name": "导出工作流", - "type": 3, - "icon": null, - "order": 6, - "permission": "module_task:workflow:export", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "导出工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "导出工作流" - }, - { - "name": "导入工作流", - "type": 3, - "icon": null, - "order": 7, - "permission": "module_task:workflow:import", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "导入工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "导入工作流" - }, { "name": "执行工作流", "type": 3, "icon": null, - "order": 8, + "order": 5, "permission": "module_task:workflow:execute", "route_name": null, "route_path": null, @@ -2981,7 +2924,7 @@ "name": "详情工作流", "type": 3, "icon": null, - "order": 9, + "order": 6, "permission": "module_task:workflow:detail", "route_name": null, "route_path": null, @@ -3000,7 +2943,7 @@ "name": "查询工作流", "type": 3, "icon": null, - "order": 10, + "order": 7, "permission": "module_task:workflow:query", "route_name": null, "route_path": null, @@ -3014,215 +2957,6 @@ "affix": false, "redirect": null, "description": "查询工作流" - }, - { - "name": "创建运行记录", - "type": 3, - "icon": null, - "order": 11, - "permission": "module_task:workflow:run:create", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "创建运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "创建运行记录" - }, - { - "name": "更新运行记录", - "type": 3, - "icon": null, - "order": 12, - "permission": "module_task:workflow:run:update", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "更新运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "更新运行记录" - }, - { - "name": "删除运行记录", - "type": 3, - "icon": null, - "order": 13, - "permission": "module_task:workflow:run:delete", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "删除运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "删除运行记录" - }, - { - "name": "清空运行记录", - "type": 3, - "icon": null, - "order": 14, - "permission": "module_task:workflow:run:delete", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "清空运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "清空运行记录" - }, - { - "name": "取消运行记录", - "type": 3, - "icon": null, - "order": 15, - "permission": "module_task:workflow:run:update", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "取消运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "取消运行记录" - }, - { - "name": "重试运行记录", - "type": 3, - "icon": null, - "order": 16, - "permission": "module_task:workflow:run:update", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "重试运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "重试运行记录" - }, - { - "name": "详情运行记录", - "type": 3, - "icon": null, - "order": 17, - "permission": "module_task:workflow:run:query", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "详情运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "详情运行记录" - }, - { - "name": "查询运行记录", - "type": 3, - "icon": null, - "order": 18, - "permission": "module_task:workflow:run:query", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "查询运行记录", - "params": null, - "affix": false, - "redirect": null, - "description": "查询运行记录" - }, - { - "name": "暂停工作流", - "type": 3, - "icon": null, - "order": 19, - "permission": "module_task:workflow:pause", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "暂停工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "暂停工作流" - }, - { - "name": "恢复工作流", - "type": 3, - "icon": null, - "order": 20, - "permission": "module_task:workflow:resume", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "恢复工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "恢复工作流" - }, - { - "name": "终止工作流", - "type": 3, - "icon": null, - "order": 21, - "permission": "module_task:workflow:terminate", - "route_name": null, - "route_path": null, - "component_path": null, - "status": "0", - "keep_alive": true, - "hidden": false, - "always_show": false, - "title": "终止工作流", - "params": null, - "affix": false, - "redirect": null, - "description": "终止工作流" } ] } @@ -3441,4 +3175,4 @@ } ] } -] +] \ No newline at end of file diff --git a/backend/app/utils/upload_util.py b/backend/app/utils/upload_util.py index cb8523c5..da48ad3a 100644 --- a/backend/app/utils/upload_util.py +++ b/backend/app/utils/upload_util.py @@ -1,6 +1,3 @@ -import hashlib -import imghdr -import mimetypes import os import random import re @@ -15,7 +12,6 @@ from app.config.setting import settings from app.core.exceptions import CustomException from app.core.logger import log - DANGEROUS_EXTENSIONS = { ".py", ".pyc", diff --git a/backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/data_level0.bin b/backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/data_level0.bin deleted file mode 100644 index dd3082eb38a7eef6982bbec3c4bf5fa5637db5fa..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 628400 zcmeIu0Sy2E0K%a6Pi+o2h(KY$fB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK zfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5 zV8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK zfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5 zV8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK zfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5 zV8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK zfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5 zV8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK zfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5 zV8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM z7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b* z1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd z0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwA zz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEj zFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r z3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@ z0|pEjFkrxd0RsjM7%*VKfB^#r3>YwAz<>b*1`HT5V8DO@0|pEjFkrxd0RsjM7%*VK Hz{|h@mp1?b diff --git a/backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/header.bin b/backend/data/chroma/ae7aebe3-8927-48ba-b976-965a1a427127/header.bin deleted file mode 100644 index 2349a18e8065afa48df1f405833655a3f39098a1..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 100 rcmZQ%K!6kk6U^$7fC#j}XsG;uC=h`16`(YX|F20q)m`+uJrTdMM3cSd@sZQH2k-Yvv3l$e!ySObw2F$W#5@kZjR8v{5v$vN&cPjZ;Ws6`{jK{ zBVUhvW~@ARa`Z>xFNbf0XF^{JDI>o z?zEf6@5+XJtD=(|pq+~7oIkmBJRCcDlI!LTQJo9U*)b0n$6uI^2ylS5Hb{eNT_~4c zDSMch0cfTM7~%3#S-8Bqun>vQJsFB!6Nd@wlzfAjJcq1?KCj^@5>p$SWR+r=QqO&yI{PWcuSf6Gsh2-?`yMr90tcgVE^yNL+%M zrR}K%UKVy1b|vv%XVubMLs+^>i^y%FR-j8P&o7o%%Cn1C>`EkDD7`qlx=sqTa1Nc$B%P6B?F`7 zZ~YB>5Kda&lMcecZKS=uNe0#;r8)S5(yZ5s(ylcdESjs5I+tiAwP{7-&p#H5T|GUl zup2VeP12AX3aQkaYqXQ_c!at~j)Y^UPjg*pjnMJ;T<-Zf;sMcf(V)RvSZ{aHu5A## zEpKen_F>CT^$8$avQuEI21Vl8!=c#xsbR%(Xy_Hm_2Wo(Vz1FeSG2=BbBDsQQ>VDs z+3@GOb)m z$BrH2y1YTz^fmuPJ%mlqUHE!RW^~anjE7&Y(OX-T<#r59(m7Y0Gry-? zj{kr{sWmpx*b?@aXN$DlL*6bSwOSpz!@d#KQ)pUvm_QA($Sf}OFLbEucMIdvf-nsQ zNeHH>aXUwtRttyd+}zLqj83uR{5$;5M&1cd1g`mcU!8hA*yD~c7LIY`jbn@Ywysvk zQg9TLX1|qms%x^*Y$qEqtv|zNlk^V^F6H%_E^NrxiEd5)^fOK*q0wv@k*JY)CgRbu z)iqYy1S~u@f3zC33lvlD`(hKHi$vxxuauU{v=uDbMeTwYqVubqor zSYEnfPSLFror}yalwn1HEp1r7EE$-a%r3&Jhm-udNGh92XR|WViX@lLr=+Tu%@A3F zUt%_w*2JCCD5z-h&TXcmQL3Us(q%^{!AAG0Wtk|=mZ}Rk*I^*HY#Gp#ZN)n#v@XbC z9f0LXmz#)sCIVWWZ41=+7Mhx1%(}pUgf$qdYNH7YaAFt_9i6B`yD7Ac?1U{@>2x++ zN%yGSl|P+U7)ocI)B`a6fuo!vc~VS^q9~`te4!|-s+P+XbE1+HrM%KD!mz|$?aGW} zm~za!t2i@gy(CxK%@$N{SDj{|I_3mKtr;^fxm6?j8IR(fD7r9v;es%?w6MB(d7#yN zE+Q&P0{SMSu@@PmbwS!h~KYK^UC zeT$fNq_GBp8|0>0N92YIW!>C{H7Te%%*Yg)(BYjC=G#UMf^*?|tzI`XHj<~w%7*zu z%5p{c* zmTddQt*w4TYGCW$bKDFeaz2wQ6w}$NCS}!3E|rl}8d0TmE+>k)Zgm9od~xS_rsol= z=WMb?Hh^5R-Q27xUd1jOR~Q0%l~z=sC1|bYhT%JBFAjh%K_x3}0mN*{dSqyoH99u4 z&KemHEVwQV2v-%;xx7-4a@B%Xm9ttgr%CxL(NslLGm@q5SDCT~IMR2&QeZ zv;T={p_;Xc2=)A zeO1aO6z$z$>Bw5p>RCmiCb(&$?G|Usl6BW|FdNS$jK|FL_`5AV*9D_NIJTba3YaOG zLIJu*QqUAoCds6vw2~La6v<@L>1)cNJXOjYlD(e&mE(6eY??}A2Co^SqhV8-kuO5P>90QeWJ~U|A=3GQ7ifIKb z#%wB`f~iC?OLdsZ6jJFtDHKKN&S{u~kCon7VVb0h6kT;{s4gAEcO77d&7&Qy^-RRw zo$3G;Ia&8=$#un=Fb-7x`@I6xuPTqE)@skWNV<^DN@`xp%BrY{xl}%*$w~(1_4^Kh z6B&zudyqfw<1_qm{;%>kc$L4#{}TT*{(tb_<^L`JFZti&|1tk1{`35&AM9EHriKKN z01`j~NB{{S0VIF~kN^@u0!ZMYA@Eo*$#M2mqU2sy!Hd^x#(LydhehADgQ(Mzw_Vczr(-H|2zKM{5Sc3!T&z20DP7I zC;UI;|33dk{_pU=$$ysrb^cTQ9at^c=G%Oe{{&C?kMkepU*@m!3;ad?JpVjj;M4pW zKF&YIpWvV1kMIZheSC=LqW>rQzoY*v`ah#Ti2j%8cccF?`tPHE82wk#KZt%K`t|7V zf*5!p0VIF~kN^@u0!RP}AOR$R1dsp{U<3kl0pD2QJauNN^NZAZfjYlHo#(0Z9Cgl7 zXNEdO>J+Gxr%sMKS?Xk{lctVD9g#XI>YSy{8R{gdlc3IN>O4!GICZ9}GewP%ASBy|MpoS@D}sB@e;Pg3U?b)KNkQR+NSofvf2~` zJWjy>3OGI={Tv`3NB{{S0VIF~kN^@u0!RP}AOR$R1Rg~KH~k~4SC6l1*h zxlBfp^4Vl2l~a=$Rm&u+g={rhl|=Z;Hc<mGli<0 ztQPaRWHFbOMfgH@9=<3(Am=UEh^g5+yWY@m!dKZ9xt@eCtLrr#Hm-UZz8}A|ycZep z`Ts|;N@JOj01`j~NB{{S0VIF~kN^@u0!RP}cnRS1|6U&W3ke_rB!C2v01`j~NB{{S z0VIF~kierzfZqQf^VNOPY^}1`Vadtg~y~Z0xB)M8AJ)6N>!C{Nr zT{jQDzr8|TG8p?o58NPxZZ7*Q+Z)yr-m@{Rn$TYVM|szoefvgJt}wQUjdPg z4Ps#jPoOQhb~vS*yVBhCwhR@6Uc9^n|5q0l=oMRT??%v6dnvkkEh(l&QIu0+zEFf6 z>$O~_m=l#8JguX2i($I$cvof|!<1v*UB#7VL#y?WT!F`VpmMwF6ARU0^k73`rCn<_ zACRK@g zhlFS|ebd}BJw#>Cxn(!2Q__0$$y*;*t}fKtuz7S7iUJK;&^w!(%~o3gLbb-4psW++ zx^T1BUKiwgo3xmk6M_tt1C*=7me%ltg_T)oT1{$=t!90Tm~{lwk%b1iY1R?C_qEd8 z+@|TnZq!CC6kr$o))`^GO*gHk=fd?`y>4b~Bu|r-g=(`|Cvrn*(=CDlJ?vH(RtwQ2S`=tX?(XF&v7$^$oWjJP)uj5nv_*Dxl~3@X+)LMxtu8Gy44ZT^TnO#Ssfjr zdd?Wqzf+Vs5p5*X@a|-G&El2klMK7GReh ztnYtIdsa*~su`QR!$#}gzRbPrv%(S0`?tXWoNotTYKSgUQpkU=+d?%y@ri^)}4jP1_Z7TQ)} z+#UM;!?IMta2XZc7|z`zg5YclZ?#bd#)`ZvU@OqBn?+lQDVCm%=jqns0D7)PWZJu(VYRd(Oi!5!fzi;A8||t1Im)n1yEAVNfE&8f z-rzzLdcGDL*;iNRFX+QEw}gY{W7QiN=&foQ7i@WG6QxITfMZ+$XDZi^9MKYOmI$IRG;Xo+HUE5|OK!8e|c8~Nn0@p||L=3mxVMntX z$+-yaizPLc&1cj!^vq&Wk*Enp_l{SJX<5E=(GUIf3vcpF6h9S3vgQUuyJvf37~UbU z=1d3o0j*gFt6|^;XD4~qz8$J8nV~j*e@|x0=nUITR?FM!#N0ixRjl2*;_GC4V2Ro- z(}tkMWx|6c?zJnJw^nb;)`jKLtYPvBZtGZ>c7$|xl$VNPS^Fx*|zOH_04(c%2E9Keq3njq`!8UC) z%`2VUpd&O)^hy^?%kH{vHgjWWRV!e^l;;;qE9Ke6D|S@^i|EDK)rGP!x4I1dZMk9x zVUeBbqLseabQoGYt(xTx&R*{j&Mz#TXX%rM(ln+}PRq7tU`Au0^;%8mvsqOG6JM&5 z0x6_ZqLxwfnxa+Ha;A{`wJkn;lzZv(l|^%>N(wX!6Z~7y*WO6h)FgBkE~xI#9s*TN zfUa2!?OeQAO)I6@fJs+dC10TbovDN#yivN@R7&z_$HwS#G) zcdGQtI=KbCBCSR7gpd@4gdoP_!gJ3FLMmk<9%jl6EChxZgxFP9mao zzlmgBiRA7#k(`r=bl+g}u0-w|Y{8YteS)=IlgrCRk=^4CzAMYaXma@h*C}dM?wrH3mutxq?)bvP4ysw3NHE9EFa5>CFkV<2M&VtN{eNeOO(D zWkvQDBmHBJr}LLDlwJ|M({X!rwaujpm5Qxh!O}thbzvHepoG!?8lmFCwM(VtlH(!F zuh;|axyZfc?KWh0%iR^u%HQSnZX|$_0oES(0Kgs4lELlg6#@qBTrLFzca{)IQBz=N zD@02Zt&qzWvZS!{{C-gC${k`V)tta-aO(t?xNfy!74JG})ILf2)%x=Ma(Q)jq1Rlh z(6}@z43lz8Sr=AU<}Y6qw6>m|cFPGfcQfhdr=aX>EqMbL4y?I0*gs3IO3<6;bgK#G z01T-vYH@2yJ>VLfBllM-&s5pDY-M*Yq80P93Px^5CbC@6VAiRnV8uRF&B05yaP7{C zamw$FxXI6)b=pZKt>)q42Wum^!p*F+{uOSUoZT_Fmj(J4xFez@XY(ly-f)IaOD-rR zpH0b6C4- zy>1%8)LyrHvIj7$u%@n#a0yq_8o0#lNNUyJ{=&Ovv377{63av$K4o?KIcFN;bp9q= zmWaD^9kRecy8*~+z?|GQ{taF*8K6F6a_3Qid#-=CF3;W9^>t6@A~0W8w5+6Ls#PhM z&k;#fbA?<=&SX+Wt&q)i(-C^7p}We~Jt9z5;X4A(v^=T9t*;Harz`2Z<;OzNc@gU# zYg!FvKJ*rz=bng@uIG*)i|xAEmsQ})1u)~U=8LM9&*h87Y9=p<1u0L|ED=?F|Nnzs z>cHHQ01`j~NB{{S0VIF~kN^@u0!RP}ydMPc{r~raGUkB5$Gx9>q4)oB z{{KiQ7Yl&|kN^@u0!RP}AOR$R1dsp{Kmtg>BycbgAM?)t+5P{o`S`CPBqV?YkN^@u z0!RP}AOR$R1dsp{KmthMVI%M&7hFxo>E|%%vD+3S)_mXECdK;Awn$<3|KIlUZ$E4$ zM0t<^5^!3BIUEmOe&`)Gpd$JRtwo`vMPziLXoJ2LcVal z(Y#qF>Kdt3Wu1h4d%0!x+U?gSO*2-XoOx}M?UPNWxuvX| zK^mQnN_)LUWEBDw5|dTAt*lq{+9x3kBMp7Kp%@|b^=2IafG4EfA~Z-vld06;c}BVo z2cX!LTRN$*(on%X9(yUCy@7m1>om;5yf)e1-lS#jG^manlaNz&L&}=!Fmi@%<7s|$ z_xSr!1dYoya*wgA&1RiOrho4(XD9^hR&UfgWgSZ9(7E5T&=l=9Py=NdJnjn%B-O69 z>j3@zYZH*?DO>$C<2^Q!p4R1UA$qgZQb+|zL;X3WHkHl>fvWPbH619sC0pfPYtAX= z9<ParWmnp84 zuCQ+~q(-HbluM?hWU3&D`57@gBjzs!?;$10N-5lvQfelXy%e~Ilte3~kT<1NVMdZJ z`R^g6l$FwmJEh!AHb0XoT;legl327-3ihOwnaPUu{{KkyFMa&C_%HB;7vU5SB!C2v z01`j~NB{{S0VIF~kN^@u0*@AfbPyIB?rp7MWbbk8RfXW*7kr`{DinKjHra|9kx3<3G#)3cto*;?w+bel+^0(eFfmKl+vE=cC=|_2}hj zA^K!AwErjjzY9_DKmter2_OL^fCP{L5#1xek{yA?HdDM=MD0j`oxhiciM=pthbsQ@{ARQrgPYg z5p>5mHUTk?Kz_lV7zB>Phs>;lRyy$6{KqB^nsP%Fdq;(X2S9Ek9Yby-Hg9GaaA){r z6f%rKhJl_8ReBuWZ{`?q<#=Rb+{_W;*f~z@gB*|2gp3@G=<4B!nWx{K=fs#9-|vcl zd=vzs@xjx>N->U!u$jEyNS=8f4w>TlUE&=+I1(O<(IOcktMr%%hQ|!0xx!5x3xvmv zqW5HI#5m-qF`zcMM4aHlV}_F5g`xNVk4ER=`+xj*`LFRm=Km4D4fFrs;>-NssaU)^2cCgP#*t291 zF(SCv3&xDBdKU{ujbzvwL72sHF$+b4V=>MUe&j%K>=fs^7CamrJI(d53mBnXV*0`3 zY1dZ*Ybw3}&+&ikgMWA+0VIF~kN^@u0!RP}AOR$R1dsp{KmtEk1g5z?ekX90?f>^S zeEgg5m=>5#c*L#wTIS`ew(m~XU6t0C7XZ`1q# zeDu$K@b>>#_&52h(Ld+^4SWZ%&A-5p!PfwOldtkg{@4doSj-g(AOR$R1dsp{Kmter z2_OL^fCP}hyCrZScxn`0m$$$Fcl2QJ6z6`SJ{CI|e0q#?e=F~?KA*$Au^v0roAAND zNPKUkBl~+JjrT>0^adI0P33T?*JmQw8wuY3e|pUG{y&}n`wqNYniwAmAOR$R1dsp{ zKmter2_OL^fCP{L5_p&i(C7d6hrjHLemnYw{XZH1#`yNWU*2~#^7Y7P#>!(SM}HLl za`;AgCiJC{GV`dewb79v807moF`?mWAb|Yx5VnBSJWK?i|-WIjsDu+?M;x zR-xXV#%NGT$t!F`wgrK)u@L)-0ORjfFJpGYS?BYqzS8J%`mVTpNYm*9n zKE?Q5*^qBlbaDd_S48Lh$*tqz*vXSzH*bjQTyV~gdAK*msb}SBJsH=L$Pb(FkzjNZ;-Z3wNnxM{hsLR$0%s-&cPngR;>Ow zM?IzKziyVwrR3$(i>2k#<+;*|UgCL1-Ht^yxp2Ht~2i%kTSGo)gk7s z)tj<)0ovgL^L5pF!B98k6_3QzPwb{=M@APi{qY^xmxe0(&J8ar-3cEXj7INA;u6d( zZBHfevaqwTD~b0ytCrpx!qQb*L~awcBDaO|{9? zuGm3XK-`~j-5Wub)hxF;L>xs=jVtAM9)Qo25Vux-9@{$LG-q~ zu}Ry9Ej!gGfMCzec{xSm*~6jO{HbBZa%kuk$@Sw%c4Dv5L|3%KJ9CG^u~VnG*Vz#3 zy5u^Y@RWYib!E4jAS9DQr%}7nA;Ma{S*1@D8a+%y-sr6yuAqQLpdDmY)`<~Z1%wTq z)(mr)Vy(dn)7J(CSKYA!8U35EbQzl1LJ6#bx!IMu*$btFpcqdtRz9IO8x`o{SkETt z$~x2s$HS1rYJODwge@X#cjCcNtaM_S&TJ_vl6`teJg~PYq8;VnPL>bHPMqLgFB&!0 zzF?mm@aW0*-les=z~7^ejJDM_$gtC06~TfGQ-g=EomxVmsN+<`QbX^wh|pe#j#CCp z3Yv9WSeJF7*{E-WLu`ErH*4*62(wn(A`K&fNIV%0#b%EU(~c=xMKn+NkL?e~jveE= zyg}LYHUC6CgiX&~_j42w+^#DYaaX71Tzxp1JMux3 zPY4Wq_nAOAHZ{fV>@#@D%8gE~MJhL2V6}UG1U-yAmrOd;3bziI-Rjait4UTE`j^gC zpud*e9o-rx?9tX}gI3^p+P|B%xlvnVqcRvt&Us+xEN57oug8tjIai!Bzo%S||A0cN zH8#-L684v8i?rNB-Yy}vS{=H>z7f?^Xj*ufKzS{V{!F1mUB6oxmllL+C`dvuMU8uM zjPDPB%opQ+O85W!Z1gXpU*G@#_J4c)YvXr9%lj^kJ|6k4@M|Mq9SH`tz&{T>HTKi7 z7e;@T`{}+P1fF`B*ECSJT?p)?&KSee>k~#3vFAv71q{HBvvyG%ek?WrJB6F6vbRxHu^$N36ZAe3|PkWFKSn`7P5?69l)7L$-6=AD)^+||_ zV@pe1_v~H7gY81}+9w?q_sM4lnA)p=R(P23*}b`|*n;@vY-;n>BCT=%1Qk#6uRrl~Es`1D?-dbjJ1PdpQj&Cl;a z<1j%R&XP_%p4zL(H1KeWu0_b5u(>$KWhW$m8IHXqa^0K8y0KGetJ=_MzP;$&e|ge` zF@O0&=@r*{+ra2_iJ?CUVH#}Yp2cn9+9gW*-d6y#?V}I&zC!3(v#;qsewN*$y6Ej}Of#?7bn$w7Z&BIR!W+o61&5dd z%nlr|iih_VEc+DiE8u+vxI_8=-dCW{|3`Sy$Nx9}AMk&TzXtF9|1tju{1^C5@WcZN zAOR$R1dsp{Kmter2_OL^fCP}hhn2woU~GbG*V=WWY1-ql;OPnYj6e%M#y10Z!#ebV zPw+qTc<}U*UVj3|L^K#X!lKik2aX0$$9lu5^oSk^#*VUp#&31s6T#C*dxI**G3tb6 z9$`Khi?PU6dW;Op#<&s~0X+u!A5v|m`FHrA@jqohKCJ4+aw7pGfCL^90>^$9-vd zM)u_Q#*cL`N+~IqOi9U9K@f8@+00B<^uvSPdQGPfRm_|{s}fD_)Z6yoLN1e0qjVjuw|fCP{L5XTX0v zn1ZJTM*RNZq0gA!&w_U}n1Y-Bfsj8C2=YJn9j1Oh&hPg@1~ik2Pnm%e;LitBaM#~I z3V{Q`L$7<2ft&oCk<184a$tmh?eC!R?mi^sga5G9;PZ{}-}Ui-!T%}$L;k>tvM*PN~Ao~+wfBgQD;Mh3lrz-*@{2%$? z$^Sp$|14+$UvB!C2v01`j~ zNB{{S0VIF~kN^@GKp+S!0#Pmq>j6=J5EcWX0eB4oUNH#58bB02q5y9e@PFome|R7P zB!C2v01`j~NB{{S0VIF~kN^@u0*?@ZeSu@#7X9jTvk|2C|2h5_eDDtsB!C2v01`j~ zNB{{S0VIF~kN^@u0!UyWf$j?j!%6P1bYEH|8&#sJwZ@uWxzQmV;!nXLB^~di;hj{P z4Xw6zM!!+7wF&(*w_KVnmxS`{`Gu0;1ZSyKSa4xFBEV5oh55_n(#6uUaAkRZad!Dt z;ib~63F8_ISowtBY*h4iODLCKDM#YxA_c9gk&LEQ(2&;JKU{*#aY9RE!8pTY?qNB{{S0VMD!5O`}TIPObxUH;_v#*cL`N+~IqgvZiS z1wqWsWHU2a(Z5OG1k>BFT0C?1tV%SwQ*YaU3%N{2k@DGOCY4i@8CA<9tA%VeS(U_M zp-9w1Azu)S6*j1(9tC1p1|)z4kN^@u0!ZMYCGhNZ_;j6N*uakVJ!XwKlgb3%y;+mD z%$nD3U-$2pnfL%Q%egY+hGv$U$)+AaW?3gQ|1MG_N{680d+6VveKmter z2_OL^fCStG-ntTUJpeH8egGgplPb(e(gSz^K*R?C+|2PD2_OL^fCP{L5_qo&oO!Pj z#&nSY5AlK?;gA3lKmter34C}6ytNp1JpwT2eFQ*E&1A9<;1K``9|8FA)F><| z5)Q{|y#N2KFZ`^NC!QezB!C2v01`j~7J;{38g)GZ@IucM0NEKa e{{Wr; jobstore?: string; executor?: string; trigger?: TriggerType; @@ -117,7 +116,6 @@ export interface NodeForm { name: string; code?: string; category?: string; - config_schema?: Record; jobstore?: string; executor?: string; func?: string; @@ -134,5 +132,7 @@ export interface NodeType { name: string; code: string; category?: string; - config_schema?: Record; + func?: string; + args?: string; + kwargs?: string; } diff --git a/frontend/src/api/module_task/workflow.ts b/frontend/src/api/module_task/workflow.ts index dfd7aa2b..3c4201cf 100644 --- a/frontend/src/api/module_task/workflow.ts +++ b/frontend/src/api/module_task/workflow.ts @@ -50,36 +50,6 @@ const WorkflowAPI = { }); }, - validateWorkflow(body: WorkflowValidateForm) { - return request>({ - url: `${API_PATH}/validate`, - method: "post", - data: body, - }); - }, - - getTemplates() { - return request>({ - url: `${API_PATH}/templates`, - method: "get", - }); - }, - - exportWorkflow(id: number) { - return request>({ - url: `${API_PATH}/export/${id}`, - method: "get", - }); - }, - - importWorkflow(body: WorkflowImportData) { - return request>({ - url: `${API_PATH}/import`, - method: "post", - data: body, - }); - }, - executeWorkflow(body: WorkflowExecuteForm) { return request>({ url: `${API_PATH}/execute`, @@ -87,130 +57,15 @@ const WorkflowAPI = { data: body, }); }, - - pauseWorkflow(id: number) { - return request({ - url: `${API_PATH}/pause/${id}`, - method: "post", - }); - }, - - resumeWorkflow(id: number) { - return request({ - url: `${API_PATH}/resume/${id}`, - method: "post", - }); - }, - - terminateWorkflow(id: number) { - return request({ - url: `${API_PATH}/terminate/${id}`, - method: "post", - }); - }, -}; - -const WorkflowRunAPI = { - getWorkflowRunList(query: WorkflowRunPageQuery) { - return request>>({ - url: `${API_PATH}/run/list`, - method: "get", - params: query, - }); - }, - - getWorkflowRunDetail(id: number) { - return request>({ - url: `${API_PATH}/run/detail/${id}`, - method: "get", - }); - }, - - createWorkflowRun(body: WorkflowRunForm) { - return request>({ - url: `${API_PATH}/run/create`, - method: "post", - data: body, - }); - }, - - updateWorkflowRun(id: number, body: WorkflowRunUpdateForm) { - return request>({ - url: `${API_PATH}/run/update/${id}`, - method: "put", - data: body, - }); - }, - - deleteWorkflowRun(ids: number[]) { - return request({ - url: `${API_PATH}/run/delete`, - method: "delete", - data: { ids }, - }); - }, - - clearWorkflowRun() { - return request({ - url: `${API_PATH}/run/clear`, - method: "delete", - }); - }, - - cancelWorkflowRun(id: number) { - return request>({ - url: `${API_PATH}/run/cancel/${id}`, - method: "post", - }); - }, - - retryWorkflowRun(id: number, retryCount: number = 1) { - return request>({ - url: `${API_PATH}/run/retry/${id}`, - method: "post", - params: { retry_count: retryCount }, - }); - }, - - pauseWorkflowRun(id: number) { - return request>({ - url: `${API_PATH}/run/pause/${id}`, - method: "post", - }); - }, - - resumeWorkflowRun(id: number) { - return request>({ - url: `${API_PATH}/run/resume/${id}`, - method: "post", - }); - }, - - terminateWorkflowRun(id: number) { - return request>({ - url: `${API_PATH}/run/terminate/${id}`, - method: "post", - }); - }, - - getWorkflowRunLogList(query: WorkflowRunLogPageQuery) { - return request>>({ - url: `${API_PATH}/run/log/list`, - method: "get", - params: query, - }); - }, }; export default WorkflowAPI; -export { WorkflowAPI, WorkflowRunAPI }; +export { WorkflowAPI }; export interface WorkflowPageQuery extends PageQuery { name?: string; code?: string; status?: string; - category?: string; - is_template?: boolean; created_time?: string[]; updated_time?: string[]; created_id?: number; @@ -224,17 +79,6 @@ export interface WorkflowTable extends BaseType { description?: string; nodes?: any[]; edges?: any[]; - version?: string; - category?: string; - tags?: string[]; - is_template?: boolean; - template_id?: number; - published_at?: string; - published_by?: number; - metadata?: Record; - canvas_state?: Record; - thumbnail?: string; - statistics?: Record; created_by?: CommonType; updated_by?: CommonType; } @@ -246,53 +90,9 @@ export interface WorkflowForm extends BaseFormType { description?: string; nodes?: any[]; edges?: any[]; - version?: string; - category?: string; - tags?: string[]; - is_template?: boolean; - template_id?: number; - metadata?: Record; } -export interface WorkflowPublishForm { - version?: string; -} - -export interface WorkflowValidateForm { - nodes: any[]; - edges: any[]; -} - -export interface WorkflowValidateResult { - is_valid: boolean; - errors: string[]; - warnings: string[]; - stats: Record; -} - -export interface WorkflowExportData { - id: number; - name: string; - code: string; - description?: string; - nodes: Record; - edges: Record; - version: string; - category?: string; - metadata?: Record; - exportedAt: string; -} - -export interface WorkflowImportData { - name?: string; - code?: string; - description?: string; - nodes: Record; - edges: Record; - version?: string; - category?: string; - metadata?: Record; -} +export interface WorkflowPublishForm {} export interface WorkflowExecuteForm { workflow_id: number; @@ -302,92 +102,11 @@ export interface WorkflowExecuteForm { } export interface WorkflowExecuteResult { - task_run_id: number; workflow_id: number; workflow_name: string; status: string; - message: string; -} - -export interface WorkflowRunPageQuery extends PageQuery { - workflow_id?: number; - workflow_name?: string; - status?: string; - business_key?: string; - initiator?: number; - job_id?: number; - created_time?: string[]; - updated_time?: string[]; -} - -export interface WorkflowRunTable extends BaseType { - workflow_id: number; - workflow_name: string; - workflow_version: string; - business_key?: string; - initiator?: number; - initiator_name?: string; + start_time?: string; + end_time?: string; variables?: Record; - status: string; - error_message?: string; - start_time?: string; - end_time?: string; - duration?: number; - retry_count: number; - max_retry: number; - job_id?: number; - metadata?: Record; - node_executions?: NodeExecution[]; -} - -export interface NodeExecution { - node_id: string; - node_name: string; - node_type: string; - status: string; - start_time?: string; - end_time?: string; - duration?: number; - error_message?: string; - input_data?: Record; - output_data?: Record; -} - -export interface WorkflowRunForm extends BaseFormType { - workflow_id: number; - workflow_name: string; - workflow_version?: string; - business_key?: string; - initiator?: number; - initiator_name?: string; - variables?: Record; - job_id?: number; -} - -export interface WorkflowRunUpdateForm extends BaseFormType { - status?: string; - error_message?: string; - start_time?: string; - end_time?: string; - duration?: number; - retry_count?: number; - max_retry_count?: number; - metadata?: Record; -} - -export interface WorkflowRunLogPageQuery extends PageQuery { - task_run_id?: number; - level?: string; - node_id?: string; - node_name?: string; - created_time?: string[]; -} - -export interface WorkflowRunLogTable extends BaseType { - task_run_id: number; - level: string; - node_id?: string; - node_name?: string; - message: string; - data?: Record; + node_results?: Record; } diff --git a/frontend/src/composables/index.ts b/frontend/src/composables/index.ts index 9dec97e5..99d083f1 100644 --- a/frontend/src/composables/index.ts +++ b/frontend/src/composables/index.ts @@ -2,3 +2,9 @@ // AI 相关 export { useAiAction } from "./ai/useAiAction"; export type { UseAiActionOptions, AiActionHandler } from "./ai/useAiAction"; + +// 任务相关 +export { useDebounce, useThrottle } from "./task/usePerformance"; +export { useNodeDrag } from "./task/useNodeDrag"; +export { useNodeOperations } from "./task/useNodeOperations"; +export { useWorkflowHistory } from "./task/useWorkflowHistory"; diff --git a/frontend/src/composables/useNodeDrag.ts b/frontend/src/composables/task/useNodeDrag.ts similarity index 90% rename from frontend/src/composables/useNodeDrag.ts rename to frontend/src/composables/task/useNodeDrag.ts index ee2f4d73..cfec5398 100644 --- a/frontend/src/composables/useNodeDrag.ts +++ b/frontend/src/composables/task/useNodeDrag.ts @@ -5,8 +5,10 @@ interface DragItem { data: { label: string; type: string; - config: any; + args?: string; + kwargs?: string; nodeId?: number; + category?: string; }; type: string; position: { x: number; y: number }; @@ -20,6 +22,8 @@ interface NodeItem { icon?: string; color?: string; class?: string; + args?: string; + kwargs?: string; } interface Coordinate { @@ -52,8 +56,10 @@ export function useNodeDrag() { data: { label: item.name, type: item.type, - config: {}, + args: item.args || "", + kwargs: item.kwargs || "{}", nodeId: item.id, + category: (item as any).category, }, type: item.type, position: { x: 0, y: 0 }, diff --git a/frontend/src/composables/useNodeOperations.ts b/frontend/src/composables/task/useNodeOperations.ts similarity index 100% rename from frontend/src/composables/useNodeOperations.ts rename to frontend/src/composables/task/useNodeOperations.ts diff --git a/frontend/src/composables/usePerformance.ts b/frontend/src/composables/task/usePerformance.ts similarity index 100% rename from frontend/src/composables/usePerformance.ts rename to frontend/src/composables/task/usePerformance.ts diff --git a/frontend/src/composables/useWorkflowHistory.ts b/frontend/src/composables/task/useWorkflowHistory.ts similarity index 100% rename from frontend/src/composables/useWorkflowHistory.ts rename to frontend/src/composables/task/useWorkflowHistory.ts diff --git a/frontend/src/layouts/components/Settings/index.vue b/frontend/src/layouts/components/Settings/index.vue index 77875449..ade0a66d 100644 --- a/frontend/src/layouts/components/Settings/index.vue +++ b/frontend/src/layouts/components/Settings/index.vue @@ -443,7 +443,6 @@ const handleCloseDrawer = () => { /* 设置抽屉样式 */ .settings-drawer { :deep(.el-drawer__body) { - position: relative; display: flex; flex-direction: column; height: 100%; @@ -461,11 +460,7 @@ const handleCloseDrawer = () => { /* 底部操作区域样式 */ .action-footer { - position: absolute; - right: 0; - bottom: 0; - left: 0; - z-index: 10; + flex-shrink: 0; padding: 0; background: var(--el-bg-color); border-top: 1px solid var(--el-border-color-light); diff --git a/frontend/src/store/modules/config.store.ts b/frontend/src/store/modules/config.store.ts index 37add232..ac7f8637 100644 --- a/frontend/src/store/modules/config.store.ts +++ b/frontend/src/store/modules/config.store.ts @@ -27,20 +27,28 @@ interface ConfigState { export const useConfigStore = defineStore("config", { state: () => ({ - configData: {} as ConfigState, // 存储系统配置 - isConfigLoaded: false, // 标记配置是否已加载 + configData: {} as ConfigState, + isConfigLoaded: false, + configLoading: false, }), actions: { async getConfig() { - const response = await ParamsAPI.getInitConfig(); - response.data.data.forEach((item: ConfigTable) => { - // 确保所有配置项都正确映射到 configData - if (item.config_value !== undefined) { - this.configData[item.config_key as keyof ConfigState] = item; - } - }); - this.isConfigLoaded = true; + if (this.isConfigLoaded || this.configLoading) { + return; + } + this.configLoading = true; + try { + const response = await ParamsAPI.getInitConfig(); + response.data.data.forEach((item: ConfigTable) => { + if (item.config_value !== undefined) { + this.configData[item.config_key as keyof ConfigState] = item; + } + }); + this.isConfigLoaded = true; + } finally { + this.configLoading = false; + } }, }, persist: true, diff --git a/frontend/src/views/module_system/auth/index.vue b/frontend/src/views/module_system/auth/index.vue index 40a81034..27494147 100644 --- a/frontend/src/views/module_system/auth/index.vue +++ b/frontend/src/views/module_system/auth/index.vue @@ -163,9 +163,6 @@ const showVoteNotification = () => { }); }; -// 组件初始化时就加载配置,而不是在onMounted中 -configStore.getConfig(); - onMounted(() => { setTimeout(showVoteNotification, 500); }); diff --git a/frontend/src/views/module_task/node/index.vue b/frontend/src/views/module_task/node/index.vue index a753cced..4405df7f 100644 --- a/frontend/src/views/module_task/node/index.vue +++ b/frontend/src/views/module_task/node/index.vue @@ -508,7 +508,6 @@ const formData = reactive({ name: "", code: undefined, category: undefined, - config_schema: undefined, jobstore: "default", executor: "default", func: defaultCodeBlock, @@ -616,8 +615,7 @@ const initialFormData: Partial = { name: "", code: undefined, category: undefined, - config_schema: undefined, - jobstore: "default", + jobstore: "sqlalchemy", executor: "default", func: defaultCodeBlock, args: undefined, diff --git a/frontend/src/views/module_task/workflow/components/DynamicNode.vue b/frontend/src/views/module_task/workflow/components/DynamicNode.vue index 1a7ab956..1c0add3c 100644 --- a/frontend/src/views/module_task/workflow/components/DynamicNode.vue +++ b/frontend/src/views/module_task/workflow/components/DynamicNode.vue @@ -7,6 +7,9 @@ >
{{ data.label }} + + {{ Object.keys(data.config).length }} +
{ - if (props.data?.type === "input") { - return { - code: "input", - name: "开始", - color: "#67c23a", - }; - } - if (props.data?.type === "output") { - return { - code: "output", - name: "结束", - color: "#f56c6c", - }; - } return { - code: "custom", - name: "自定义节点", - color: "#409EFF", + code: props.data?.type || "custom", + name: props.data?.label || "自定义节点", + color: getCategoryColor(props.data?.category), }; }); +function getCategoryColor(category) { + const colorMap = { + trigger: "#e6a23c", + action: "#409eff", + condition: "#67c23a", + control: "#909399", + }; + return colorMap[category] || "#409eff"; +} + const nodeClass = computed(() => { if (props.data?.type === "input") { return "start-node"; @@ -153,6 +152,15 @@ const nodeClass = computed(() => { letter-spacing: 0.5px; } +.node-badge { + padding: 0 6px; + font-size: 10px; + font-weight: 500; + color: #fff; + background: #409eff; + border-radius: 10px; +} + .vue-flow__handle { opacity: 0; transition: opacity 0.2s ease; diff --git a/frontend/src/views/module_task/workflow/components/NodeConfigPanel.vue b/frontend/src/views/module_task/workflow/components/NodeConfigPanel.vue index a71ad515..b8d1c140 100644 --- a/frontend/src/views/module_task/workflow/components/NodeConfigPanel.vue +++ b/frontend/src/views/module_task/workflow/components/NodeConfigPanel.vue @@ -24,37 +24,19 @@ - -
- - - - - - - - -
+ + +
多个参数用逗号分隔
+
+ + + +
JSON 格式的关键字参数
@@ -84,8 +66,6 @@ import { ElInput, ElSelect, ElOption, - ElInputNumber, - ElSwitch, ElMessage, ElIcon, } from "element-plus"; @@ -102,13 +82,12 @@ const props = defineProps({ const emit = defineEmits(["close", "save", "delete"]); const nodeTypes = ref([]); -const selectedTemplate = ref(); -const nodeConfigSchema = ref([]); const formData = ref({ type: props.node?.type || "", label: props.node?.data?.label || "", - config: props.node?.data?.config || {}, + args: props.node?.data?.args || "", + kwargsStr: props.node?.data?.kwargsStr || "{}", description: props.node?.data?.description || "", }); @@ -125,33 +104,36 @@ const loadNodeTypes = async () => { const handleTypeChange = async (typeCode: string) => { const nodeType = nodeTypes.value.find((t) => t.code === typeCode); - if (nodeType && nodeType.config_schema) { - nodeConfigSchema.value = nodeType.config_schema.fields || []; - } else { - nodeConfigSchema.value = []; + if (nodeType) { + formData.value.args = nodeType.args || ""; + formData.value.kwargsStr = nodeType.kwargs || "{}"; } - - selectedTemplate.value = undefined; - formData.value.config = {}; }; watch( () => props.node, (newNode) => { if (newNode) { + const kwargsData = newNode.data?.kwargs; + let kwargsStr = "{}"; + if (kwargsData) { + if (typeof kwargsData === "string") { + kwargsStr = kwargsData; + } else if (typeof kwargsData === "object") { + kwargsStr = JSON.stringify(kwargsData, null, 2); + } + } + formData.value = { type: newNode.type || "", label: newNode.data?.label || "", - config: newNode.data?.config || {}, + args: newNode.data?.args || "", + kwargsStr, description: newNode.data?.description || "", }; - - if (newNode.type) { - handleTypeChange(newNode.type); - } } }, - { deep: true } + { deep: true, immediate: true } ); function handleClose() { @@ -159,7 +141,22 @@ function handleClose() { } function handleSave() { - emit("save", formData.value); + try { + if (formData.value.kwargsStr && formData.value.kwargsStr.trim()) { + JSON.parse(formData.value.kwargsStr); + } + } catch { + ElMessage.error("关键字参数 JSON 格式错误"); + return; + } + + emit("save", { + type: formData.value.type, + label: formData.value.label, + args: formData.value.args, + kwargs: formData.value.kwargsStr, + description: formData.value.description, + }); ElMessage.success("保存成功"); } @@ -201,8 +198,10 @@ onMounted(() => { overflow-y: auto; } -.config-field { - margin-bottom: 8px; +.field-hint { + margin-top: 4px; + font-size: 12px; + color: #909399; } .panel-actions { diff --git a/frontend/src/views/module_task/workflow/components/WorkflowDesignDrawer.vue b/frontend/src/views/module_task/workflow/components/WorkflowDesignDrawer.vue index d2e1881b..b2282d27 100644 --- a/frontend/src/views/module_task/workflow/components/WorkflowDesignDrawer.vue +++ b/frontend/src/views/module_task/workflow/components/WorkflowDesignDrawer.vue @@ -26,18 +26,6 @@ - - - - - - - - +