工作流调度引擎实现计划
需求概述
实现智能体编排后的工作流执行调度引擎。核心要求:
- 融入 Skill 思想 — 单个节点(尤其是 Agent 节点)执行时参考对应 Skill 的 inputs/outputs 定义,灵活处理
- 总体遵循格式化输入输出 — 工作流整体有明确的输入(用户输入节点)和输出(输出节点),数据在节点间通过变量名传递
- 确保不丢节点 — 条件分支后必须汇聚,DAG 拓扑排序保证所有可达节点都被执行
- 不要太死板 — 模板渲染容错,变量缺失时优雅降级而非报错中断
现有架构
graphData 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<String, Object> // 变量存储池,节点间数据传递
- nodeId: String // 当前执行节点 ID
- workflowId: Long // 工作流 ID
NodeExecutionResult(节点执行结果)
- nodeId: String
- status: RUNNING | SUCCESS | FAILED | SKIPPED
- output: Map<String, Object> // 节点输出,写入 ExecutionContext
- error: String // 失败时的错误信息
1.2 NodeExecutor 接口与实现
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 |
支持状态样式 |
执行顺序
- ✅ 输出计划文档
- Phase 1 — 后端引擎核心
- Phase 2 — 前端执行面板
- Phase 3 — 模板渲染与调试