workflow-scheduling-engine-plan.md 7.7 KB

工作流调度引擎实现计划

需求概述

实现智能体编排后的工作流执行调度引擎。核心要求:

  1. 融入 Skill 思想 — 单个节点(尤其是 Agent 节点)执行时参考对应 Skill 的 inputs/outputs 定义,灵活处理
  2. 总体遵循格式化输入输出 — 工作流整体有明确的输入(用户输入节点)和输出(输出节点),数据在节点间通过变量名传递
  3. 确保不丢节点 — 条件分支后必须汇聚,DAG 拓扑排序保证所有可达节点都被执行
  4. 不要太死板 — 模板渲染容错,变量缺失时优雅降级而非报错中断

现有架构

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.2(智谱 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 支持状态样式

执行顺序

  1. ✅ 输出计划文档
  2. Phase 1 — 后端引擎核心
  3. Phase 2 — 前端执行面板
  4. Phase 3 — 模板渲染与调试