workflow-variable-scope-design.md 13 KB

工作流节点变量作用域与输入变量关联设计

1. 背景与目标

当前工作流引擎使用扁平变量表 WorkflowContext.variables 传递节点输出:每个节点执行完成后,把输出 Map 中的 key-value 直接写入上下文,下游节点通过 {{变量名}} 模板或 ContextPromptHelper 自动注入来消费变量。

该机制存在两个问题:

  1. 同名变量覆盖:若多个前驱节点都输出同名变量(例如都叫 result),后执行的节点会覆盖前者,导致下游节点只能看到最近一个同名输出。
  2. 无法精确关联:节点输入变量无法显式指定来源;隐式关联仅依赖同名匹配,不支持类型转换、JSON 嵌套路径、模糊匹配等场景。

本设计目标:

  • 为每个节点构建隔离的节点级工作区,仅包含当前节点可达前驱节点的输出变量。
  • 所有前驱节点输出按 节点ID.变量名 命名空间化,避免同名覆盖。
  • 支持输入变量显式关联前驱输出,并自动按类型转换规则转换。
  • 支持四级隐式关联 fallback:精确同名 → 忽略分隔符/大小写同名 → JSON 嵌套路径同名 → JSON 嵌套路径模糊同名。

2. 核心概念

2.1 节点工作区(NodeWorkspace)

节点 N 开始执行前,引擎为其构建一个工作区:

NodeWorkspace(N) = {
  variables: Map<String, Object>,        // 可直接通过 {{var}} 引用的扁平变量
  scopedOutputs: Map<String, Map<String, Object>>,  // nodeId -> {varName -> value}
  files: List<String>,                   // 工作空间文件列表
  runId: String
}

其中 variables 的构造规则:

  • 保留初始运行输入(request.inputs)。
  • N 的每个可达前驱节点 P,把 scopedOutputs[P] 中的变量以 P_varName 形式注入(双下划线或点号分隔,见下文)。
  • 是否同时保留扁平变量(无节点前缀)作为向后兼容,通过 feature flag 控制。

2.2 节点级命名空间

WorkflowContext 中新增独立的节点输出存储:

private final Map<String, Map<String, Object>> nodeScopedOutputs = new ConcurrentHashMap<>();

每个节点执行成功后:

context.setNodeOutput(nodeId, output);        // 现有:扁平合并(保留兼容)
context.setNodeScopedOutput(nodeId, output);  // 新增:命名空间隔离

2.3 变量引用语法

为兼容现有 {{varName}} 模板,同时支持精确引用某个节点的输出,引入两种引用方式:

语法 含义 示例
{{varName}} 引用当前工作区中的 varName(可能是扁平变量,也可能是最近同名变量) {{question}}
{{nodeId.varName}} 引用指定节点 nodeId 的输出变量 varName {{node_2.result}}
{{nodeId.var.path}} 引用指定节点输出变量的嵌套 JSON 路径 {{node_2.output.total.result}}

注:nodeId 本身可能包含下划线,为消除歧义,内部实现使用 __NODE__ 作为分隔符或采用 Map 结构,前端展示使用 . 语法。


3. 数据模型

3.1 后端 WorkflowContext 扩展

public class WorkflowContext {
    private final Map<String, Object> variables = new ConcurrentHashMap<>();
    private final Map<String, Map<String, Object>> nodeScopedOutputs = new ConcurrentHashMap<>();
    private final List<NodeOutput> nodeOutputs = ...;
    private final Path workingDir;
    // ...

    public void setNodeScopedOutput(String nodeId, Map<String, Object> output);
    public Map<String, Object> getNodeScopedOutput(String nodeId);
    public Map<String, Map<String, Object>> getAllNodeScopedOutputs();
}

3.2 前端节点 data 扩展

节点 data.inputs[] 中每个字段增加 mapping

{
  "name": "result",
  "type": "string",
  "description": "",
  "required": true,
  "mapping": {
    "sourceNodeId": "node_2",
    "sourceField": "result",
    "sourcePath": null,
    "autoMapped": false
  }
}

sourcePath 非空时,表示从 sourceField 的 JSON 嵌套路径取值,例如 output.total.result

autoMapped 标记该关联是否由隐式关联算法自动产生,用户可覆盖。

3.3 边 mapping 数据

保留现有边 data.mapping 用于精确源→目标字段映射,同时新增 autoResolved 标记:

{
  "mapping": [
    {"sourceField": "result", "targetField": "result"}
  ],
  "autoResolved": true
}

4. 算法设计

4.1 可达前驱计算

基于现有 DagResolver.ResolvedDag 增加反向邻接表:

Map<String, List<EdgeInfo>> incomingEdges

静态可达前驱(规划阶段使用):

Set<String> getReachablePredecessors(String nodeId) {
    Set<String> visited = new HashSet<>();
    Queue<String> queue = new LinkedList<>();
    queue.add(nodeId);
    while (!queue.isEmpty()) {
        String cur = queue.poll();
        for (EdgeInfo edge : incomingEdges.getOrDefault(cur, emptyList())) {
            if (visited.add(edge.getSource())) {
                queue.add(edge.getSource());
            }
        }
    }
    visited.remove(nodeId);
    return visited;
}

动态可达前驱(运行时,条件分支生效后):

  • 条件节点仅把实际激活的分支下游加入可达集。
  • 实现方式:WorkflowLevelExecutor 在节点执行后根据 selectedBranch 维护 activeEdges,动态可达集基于活跃边计算。

4.2 节点工作区构建

NodeWorkspace buildWorkspace(String currentNodeId, ResolvedDag dag, WorkflowContext context) {
    NodeWorkspace ws = new NodeWorkspace();
    ws.getVariables().putAll(context.getInitialInputs());

    Set<String> reachable = dag.getReachablePredecessors(currentNodeId);
    for (String nodeId : topologicalOrder(reachable)) {
        Map<String, Object> scoped = context.getNodeScopedOutput(nodeId);
        if (scoped == null) continue;
        ws.getScopedOutputs().put(nodeId, scoped);
        for (Map.Entry<String, Object> e : scoped.entrySet()) {
            String scopedKey = nodeId + "__" + e.getKey();
            ws.getVariables().put(scopedKey, e.getValue());
            // 向后兼容:同时扁平注入(可选,受 feature flag 控制)
            ws.getVariables().putIfAbsent(e.getKey(), e.getValue());
        }
    }

    ws.setFiles(scanWorkingDir(context.getWorkingDir()));
    return ws;
}

4.3 输入变量解析算法

对节点 N 的每个输入字段 input

1. 若 input.mapping 存在且 sourceNodeId 属于 reachable(N):
    value = resolvePath(scopedOutputs[sourceNodeId], sourceField, sourcePath)
    return convert(value, input.type)

2. 否则进入隐式关联:
    candidates = 按拓扑逆序排列的 reachable(N) 中的前驱节点

    2.1 精确同名:在 candidates 中查找第一个输出字段名 == input.name
    2.2 模糊同名:忽略 [-_] 和大小写后匹配
    2.3 JSON 嵌套精确同名:遍历 candidates 的所有输出变量,按 JSON Path 查找字段名 == input.name
    2.4 JSON 嵌套模糊同名:忽略 [-_] 和大小写后匹配

3. 若找到匹配:
    记录 mapping = {sourceNodeId, sourceField, sourcePath, autoMapped: true}
    return convert(value, input.type)

4. 若未找到:
    若 input.required 为 true,前置条件校验失败
    否则返回 null

4.4 模糊匹配归一化

static String normalizeForMatch(String s) {
    return s.replaceAll("[-_]", "").toLowerCase(Locale.ROOT);
}

4.5 JSON Path 遍历

支持点号路径,数组使用数字下标:

  • output.total.resultvalue["output"]["total"]["result"]
  • items.0.namevalue["items"][0]["name"]

实现基于 Jackson JsonNode 或递归 Map,缺失路径返回 null


5. 类型转换规则

基于 docs/workflow-node-fields.md 实现 VariableConverter

5.1 通用转换

源类型 目标 string 目标 number 目标 boolean 目标 object 目标 array
String 原值 parseDouble / NaN true/false/1/0/yes/no JSON parse / wrap JSON parse / wrap
Number toString 原值 != 0 wrap wrap
Boolean toString 1/0 原值 wrap wrap
Map JSON string NaN !empty 原值 entry list
Collection JSON string size !empty firstOrWrap 原值

5.2 节点特定转换

KnowledgeRetrieval.evidences(List)
目标类型 行为
string 取 Top-1 的 content(可配置 top1/concat/json)
number evidenceCount
boolean 非空
object Top-1 证据 Map
array Top-K 证据 Map 列表

LLM/Agent/SmartAction.result(String)

目标类型 行为
string 原值
number 提取首个数字 / NaN
boolean true/false/1/0/yes/no / 非空
object JSON parse / wrap
array JSON parse / lines split

6. 接口设计

6.1 新增 Java 类

类名 职责
NodeWorkspace 节点工作区数据对象
NodeWorkspaceBuilder 根据 DAG 和上下文构建工作区
NodeInputResolver 解析节点输入变量,支持显式/隐式关联
VariableConverter 类型转换器
JsonPathExtractor JSON 嵌套路径取值
WorkflowScopeResolver 可达前驱计算(静态+动态)

6.2 修改的 Java 类

类名 修改点
WorkflowContext 增加 nodeScopedOutputs 及相关方法
DagResolver.ResolvedDag 增加 incomingEdges
WorkflowLevelExecutor 节点执行前构建工作区、解析输入;条件分支维护活跃边
TemplateRenderer 支持 {{nodeId.varName}}{{nodeId.var.path}} 语法
各 Executor 从 NodeWorkspace 读取输入,不再直接读扁平 variables

6.3 前端改动

文件 修改点
ioInference.js getNodeInputs 返回 mapping 字段;inferEdgeMapping 增加隐式关联逻辑
WorkflowEditor.vue 边变化时刷新相关节点 input mapping
节点配置面板 每个 input 增加“关联变量”选择器
workflowNode.js 默认 data 中 inputs 增加 mapping 字段

7. 兼容策略

7.1 向后兼容

  • 保留 WorkflowContext.variables 扁平表,旧工作流不开启命名空间也能运行。
  • 新增 feature flag workflow.node-scoped-variables.enabled,默认 false
  • 当 flag 关闭时:
    • 仍然写入 nodeScopedOutputs(无影响)。
    • NodeWorkspaceBuilder 仍然构造工作区,但 variables 同时保留扁平注入。
    • 各执行器优先从工作区读取,未找到时 fallback 到扁平 variables。

7.2 数据迁移

  • 旧 graphData 中的 inputs 没有 mapping 字段,解析时按 null 处理,走隐式关联。
  • 前端保存时自动补全 mapping 字段。

7.3 条件分支

  • 静态规划阶段使用保守可达集(所有分支都视为可达)。
  • 运行时根据 selectedBranch 动态过滤,确保未激活分支的变量不进入下游工作区。

8. 测试策略

8.1 单元测试

  • WorkflowScopeResolverTest:线性 DAG、分支 DAG、循环检测。
  • NodeInputResolverTest:四级隐式关联 fallback、显式 mapping、关联失效。
  • VariableConverterTest:通用转换表、RagEvidence 转换、LLM result 转换。
  • JsonPathExtractorTest:对象嵌套、数组下标、缺失路径。

8.2 集成测试

  • 复杂 DAG 执行后,验证下游节点工作区只包含可达前驱变量。
  • 同名变量场景:多个前驱输出同名变量,下游通过 nodeId.varName 精确引用。

8.3 回归测试

  • 旧工作流(无 mapping)执行结果与改造前一致。

9. 实现顺序建议

为控制风险,建议按以下 MVP → 增强 → 完善的顺序实现:

  1. MVP

    • WorkflowContext 节点级输出存储
    • DagResolver 反向索引
    • NodeWorkspace 构建
    • VariableConverter 通用转换
    • NodeInputResolver 显式 mapping + 精确同名隐式关联
    • TemplateRenderer 支持 {{nodeId.varName}}
    • 修改 WorkflowLevelExecutor 和主要 Executor
  2. 增强

    • 2.2 / 2.3 / 2.4 级隐式关联
    • 动态条件分支可达集
  3. 完善

    • 前端 mapping UI
    • 自动保存 mapping
    • feature flag 默认开启并移除兼容代码

10. 相关文档

  • docs/workflow-node-fields.md:节点自动输入输出字段与转换规则
  • docs/workflow-scheduling-engine-plan.md:工作流调度引擎设计

TODO

☐ 将设计写入 docs/workflow-variable-scope-design.md ☐ 阶段1:后端基础设施(WorkflowContext、DagResolver、NodeWorkspace、VariableConverter) ☐ 阶段2:输入变量解析(NodeInputResolver、WorkflowLevelExecutor、各执行器) ☐ 创建 NodeInputResolver(显式 mapping + 四级隐式 fallback) ☐ 更新 TemplateRenderer 支持 {{nodeId.varName}} 语法 ☐ 更新 ContextPromptHelper 支持 NodeWorkspace / Map 变量 ☐ 更新 NodeExecutor 接口与各执行器使用 NodeWorkspace ☐ 更新 NodeWorkspaceBuilder 为 Spring Bean 并支持 activeEdges ☐ 更新 WorkflowLevelExecutor 构建工作区、解析输入、维护活跃边 ☐ 编译并修复错误 ☐ 阶段3:前端配置界面(data结构、ioInference、mapping UI) ☐ 阶段4:兼容性与测试 ☐ 添加 NodeInputResolver / TemplateRenderer / VariableConverter / WorkflowScopeResolver 测试