# 工作流调度引擎实现计划 ## 需求概述 实现智能体编排后的工作流执行调度引擎。核心要求: 1. **融入 Skill 思想** — 单个节点(尤其是 Agent 节点)执行时参考对应 Skill 的 inputs/outputs 定义,灵活处理 2. **总体遵循格式化输入输出** — 工作流整体有明确的输入(用户输入节点)和输出(输出节点),数据在节点间通过变量名传递 3. **确保不丢节点** — 条件分支后必须汇聚,DAG 拓扑排序保证所有可达节点都被执行 4. **不要太死板** — 模板渲染容错,变量缺失时优雅降级而非报错中断 ## 现有架构 ### graphData JSON 结构 ```json { "nodes": [ {"id": "node-xxx", "type": "userInput|llm|agent|condition|output", "position": {...}, "data": {...}} ], "edges": [ {"id": "edge-xxx", "source": "node-1", "target": "node-2", "sourceHandle": "output|branch-0|branch-1"} ] } ``` ### 节点类型 data 结构 | 类型 | data 字段 | 说明 | |------|----------|------| | userInput | `variables: [{name, label, type, description}]` | 定义工作流输入变量 | | llm | `model, systemPrompt, userPrompt` | 大模型调用,支持 `{{变量名}}` 模板 | | agent | `agentId(folderName), agentName` | 关联 Skill,读取 SKILL.md 内容作为系统提示 | | condition | `conditions: [{type: IF/ELIF, expression}]` + 隐含 ELSE | 条件分支,由 LLM 判断走哪条分支 | | output | `outputType: text/json` | 工作流最终输出 | ### 条件分支边关系 - `sourceHandle: "branch-0"` → IF 分支 - `sourceHandle: "branch-1"` → ELIF 分支 - `sourceHandle: "branch-2"` → ELSE 分支 ### 已有基础设施 - **Spring AI ChatClient** — 已配置 GLM-5.1(智谱 AI) - **SSE 推送** — SseService 已实现,可复用 - **Skill 系统** — SkillDTO 含 inputs/outputs,SkillMarkdownParser 可读取 SKILL.md - **IOField** — 输入输出字段定义模型 --- ## Phase 1: 后端执行引擎核心 ### 1.1 核心模型 #### ExecutionContext(执行上下文) ``` - variables: Map // 变量存储池,节点间数据传递 - nodeId: String // 当前执行节点 ID - workflowId: Long // 工作流 ID ``` #### NodeExecutionResult(节点执行结果) ``` - nodeId: String - status: RUNNING | SUCCESS | FAILED | SKIPPED - output: Map // 节点输出,写入 ExecutionContext - error: String // 失败时的错误信息 ``` ### 1.2 NodeExecutor 接口与实现 ```java public interface NodeExecutor { String getType(); // userInput | llm | agent | condition | output NodeExecutionResult execute(JsonNode data, ExecutionContext context); } ``` **5 个实现类:** | 执行器 | 核心逻辑 | |--------|---------| | UserInputExecutor | 从 context.variables 中提取用户输入值,放入上下文 | | LlmExecutor | 渲染 systemPrompt/userPrompt 中的 `{{变量名}}`,调用 ChatClient,结果写入上下文 | | AgentExecutor | 根据 agentId 读取 SKILL.md 内容作为系统提示,结合 Skill 的 inputs 从上下文取值构造用户消息,调用 ChatClient | | ConditionExecutor | 用 LLM 判断走哪个分支(根据 expression 和上下文数据),返回选中的 sourceHandle | | OutputExecutor | 从上下文收集输出,格式化为最终结果 | ### 1.3 DagResolver(图解析与调度) ``` 输入: graphData JSON 输出: 按拓扑排序的执行层级 核心逻辑: 1. 解析 nodes 和 edges 2. 构建邻接表(按 sourceHandle 分组) 3. Kahn 算法拓扑排序 4. 处理条件分支: 记录每个 condition 节点的分支映射(sourceHandle → 目标节点列表) 5. 运行时: 按层级执行,condition 节点返回选中分支后,动态标记活跃/跳过路径 ``` ### 1.4 WorkflowEngine(调度引擎) ``` execute(workflowId, userInput): 1. 加载 Workflow,解析 graphData 2. DagResolver 解析图结构,得到拓扑层级 3. 初始化 ExecutionContext,填充用户输入 4. 按层级遍历: a. 找到该层所有活跃节点 b. 并行执行活跃节点(CompletableFuture) c. 条件节点特殊处理: 返回选中的 sourceHandle,标记后续不活跃分支的节点为 SKIPPED d. 每个节点完成后通过 SSE 推送状态 5. 收集 Output 节点结果,返回最终输出 ``` ### 1.5 执行 API ``` POST /api/workflows/{id}/run Body: { "inputs": { "变量名": "值" } } Response: SSE stream - event: node_status data: { nodeId, status, output? } - event: workflow_complete data: { outputs: [...], success: true/false } ``` ### 1.6 SSE 扩展 扩展现有 SseService,新增工作流执行事件推送方法。 --- ## Phase 2: 前端执行面板 ### 2.1 编辑器增强 - 顶部工具栏增加「运行」按钮 - 点击运行 → 弹出输入对话框(根据 userInput 节点的 variables 定义动态生成表单) - 执行过程中节点实时高亮状态(待执行/执行中/已完成/失败/跳过) - 结果面板展示最终输出 ### 2.2 SSE 连接 - 工作流编辑器页面建立 SSE 连接 - 接收 node_status 事件,更新对应节点的视觉状态 - 接收 workflow_complete 事件,展示最终结果 ### 2.3 节点状态样式 | 状态 | 样式 | |------|------| | pending | 默认 | | running | 脉冲动画 + 蓝色边框 | | success | 绿色边框 | | failed | 红色边框 | | skipped | 灰色半透明 | --- ## Phase 3: 模板渲染与调试优化 ### 3.1 模板渲染 `{{变量名}}` 语法在 systemPrompt 和 userPrompt 中使用: ``` 输入: "请分析以下内容: {{content}}" 上下文: { content: "今天天气不错" } 输出: "请分析以下内容: 今天天气不错" ``` 容错策略: 变量不存在时保留原始 `{{变量名}}`,不报错。 ### 3.2 调试支持 - 每个节点执行完成后记录输入/输出日志 - 执行失败时返回详细错误信息和执行路径 - 支持查看单次执行的完整日志 --- ## 文件清单(新增/修改) ### 新增后端文件 | 文件 | 说明 | |------|------| | `model/dto/ExecutionContext.java` | 执行上下文 | | `model/dto/NodeExecutionResult.java` | 节点执行结果 | | `model/dto/WorkflowRunRequest.java` | 运行请求 DTO | | `model/dto/WorkflowRunEvent.java` | SSE 事件 DTO | | `engine/NodeExecutor.java` | 节点执行器接口 | | `engine/DagResolver.java` | DAG 图解析与拓扑排序 | | `engine/WorkflowEngine.java` | 工作流调度引擎 | | `engine/TemplateRenderer.java` | `{{变量名}}` 模板渲染 | | `engine/executor/UserInputExecutor.java` | 用户输入执行器 | | `engine/executor/LlmExecutor.java` | LLM 执行器 | | `engine/executor/AgentExecutor.java` | Agent 执行器 | | `engine/executor/ConditionExecutor.java` | 条件分支执行器 | | `engine/executor/OutputExecutor.java` | 输出执行器 | ### 修改后端文件 | 文件 | 修改内容 | |------|---------| | `controller/WorkflowController.java` | 新增 `/run` 端点 | | `service/WorkflowService.java` | 新增执行方法签名 | | `service/impl/WorkflowServiceImpl.java` | 委托 WorkflowEngine | | `service/SseService.java` | 新增工作流事件推送方法 | ### 新增前端文件 | 文件 | 说明 | |------|------| | `components/workflow/RunDialog.vue` | 运行输入对话框 | | `components/workflow/ResultPanel.vue` | 执行结果面板 | ### 修改前端文件 | 文件 | 修改内容 | |------|---------| | `views/workflow/WorkflowEditor.vue` | 运行按钮、SSE 连接、节点状态样式 | | `api/workflow.js` | 新增 runWorkflow API | | `components/workflow/nodes/BaseNode.vue` | 支持状态样式 | --- ## 执行顺序 1. ✅ 输出计划文档 2. Phase 1 — 后端引擎核心 3. Phase 2 — 前端执行面板 4. Phase 3 — 模板渲染与调试