Răsfoiți Sursa

1. 重构了工作流变量作用域隔离机制,移除 WorkflowContext 中的扁平变量表与 scopeMode 兼容模式开关,所有执行路径统一采用节点命名空间({nodeId}__{varName})严格隔离,杜绝不同分支同名输出互相覆盖。

weisijie 1 lună în urmă
părinte
comite
15f2526e91

+ 0 - 47
backend/src/main/java/com/agent/management/engine/ContextPromptHelper.java

@@ -52,29 +52,6 @@ public final class ContextPromptHelper {
         return buildContextSystemPrompt(inputValues, "节点输入变量");
     }
 
-    /**
-     * 把工作流上下文的 variables 格式化为系统提示词片段。
-     *
-     * @param context 工作流上下文(可能为 null,调用方方便起见)
-     * @return 格式化后的提示词;如果上下文为空,返回空串
-     */
-    public static String buildContextSystemPrompt(WorkflowContext context) {
-        if (context == null) {
-            return "";
-        }
-        return buildContextSystemPrompt(context.getVariables());
-    }
-
-    /**
-     * 把工作区变量表格式化为系统提示词片段。
-     *
-     * @param variables 变量表(可能为 null)
-     * @return 格式化后的提示词;如果变量表为空,返回空串
-     */
-    public static String buildContextSystemPrompt(Map<String, Object> variables) {
-        return buildContextSystemPrompt(variables, "工作流上下文(来自前置节点的输出变量)");
-    }
-
     /**
      * 把变量表格式化为系统提示词片段。
      *
@@ -99,30 +76,6 @@ public final class ContextPromptHelper {
         return sb.toString();
     }
 
-    /**
-     * 把上下文系统提示词合并到用户自定义的系统提示词后。
-     * 两者皆空时返回空串;其中一方为空时返回另一方。
-     *
-     * @param userSystemPrompt 节点本身声明的 systemPrompt(已渲染)
-     * @param context          工作流上下文
-     * @return 合并后的完整 systemPrompt
-     */
-    public static String merge(String userSystemPrompt, WorkflowContext context) {
-        return merge(userSystemPrompt, context == null ? null : context.getVariables());
-    }
-
-    /**
-     * 把上下文系统提示词合并到用户自定义的系统提示词后。
-     *
-     * @param userSystemPrompt 节点本身声明的 systemPrompt(已渲染)
-     * @param variables        变量表
-     * @return 合并后的完整 systemPrompt
-     */
-    public static String merge(String userSystemPrompt, Map<String, Object> variables) {
-        String contextPrompt = buildContextSystemPrompt(variables);
-        return join(userSystemPrompt, contextPrompt);
-    }
-
     /**
      * 把节点明确定义的输入变量作为上下文系统提示词,合并到用户自定义的系统提示词后。
      *

+ 1 - 1
backend/src/main/java/com/agent/management/engine/NodeInputResolver.java

@@ -108,7 +108,7 @@ public class NodeInputResolver {
             }
         }
 
-        // 2. 直接变量表(覆盖初始输入与向后兼容的扁平变量
+        // 2. 初始输入(运行级变量,扁平形式,不涉及节点命名空间
         Object direct = workspace.getVariable(name);
         if (direct != null) {
             return VariableConverter.convert(direct, type);

+ 13 - 24
backend/src/main/java/com/agent/management/engine/NodeWorkspaceBuilder.java

@@ -11,6 +11,14 @@ import java.util.stream.Stream;
 
 /**
  * 根据 DAG 与 WorkflowContext 构建节点工作区。
+ *
+ * <p>变量注入策略:严格作用域隔离 ——
+ * <ul>
+ *   <li>初始输入以扁平形式注入(运行级 key,不涉及节点冲突)</li>
+ *   <li>每个可达前驱的输出以 {predecessorId}__{varName} 形式注入,互不覆盖</li>
+ *   <li>不再注入无前缀的扁平形式,避免同名输出互相覆盖</li>
+ * </ul>
+ * 后继节点通过 {@link NodeInputResolver} 的显式 mapping 或四级隐式回退取值。
  */
 @Component
 public class NodeWorkspaceBuilder {
@@ -26,15 +34,9 @@ public class NodeWorkspaceBuilder {
 
     /**
      * 为指定节点构建工作区(静态可达前驱)。
-     *
-     * @param nodeId      当前节点 ID
-     * @param dag         已解析的 DAG
-     * @param context     工作流上下文
-     * @param scopeMode   是否启用严格作用域隔离(true:只注入可达前驱变量;false:额外保留扁平变量作为兼容)
      */
-    public NodeWorkspace build(String nodeId, DagResolver.ResolvedDag dag,
-                                WorkflowContext context, boolean scopeMode) {
-        return build(nodeId, dag, context, Collections.emptySet(), scopeMode);
+    public NodeWorkspace build(String nodeId, DagResolver.ResolvedDag dag, WorkflowContext context) {
+        return build(nodeId, dag, context, Collections.emptySet());
     }
 
     /**
@@ -43,16 +45,14 @@ public class NodeWorkspaceBuilder {
      * @param nodeId      当前节点 ID
      * @param dag         已解析的 DAG
      * @param context     工作流上下文
-     * @param activeEdges 活跃边 key 集合(source-sourceHandle-target);null 或空表示使用静态可达集
-     * @param scopeMode   是否启用严格作用域隔离
+     * @param activeEdges 活跃边 key 集合(source-sourceHandle-target);空集合表示使用静态可达集
      */
     public NodeWorkspace build(String nodeId, DagResolver.ResolvedDag dag,
-                                WorkflowContext context, Set<String> activeEdges,
-                                boolean scopeMode) {
+                                WorkflowContext context, Set<String> activeEdges) {
         NodeWorkspace workspace = new NodeWorkspace();
         workspace.setRunId(context.getRunId());
 
-        // 1. 初始输入始终可用
+        // 1. 初始输入始终可用(扁平形式)
         workspace.getVariables().putAll(context.getInitialInputs());
 
         // 2. 收集可达前驱节点,按拓扑顺序注入命名空间变量
@@ -69,10 +69,6 @@ public class NodeWorkspaceBuilder {
             for (Map.Entry<String, Object> e : scoped.entrySet()) {
                 String scopedKey = predecessorId + NODE_FIELD_SEPARATOR + e.getKey();
                 workspace.getVariables().put(scopedKey, e.getValue());
-                // 不启用严格隔离时,同时保留扁平变量(最近前驱覆盖)
-                if (!scopeMode) {
-                    workspace.getVariables().put(e.getKey(), e.getValue());
-                }
             }
         }
 
@@ -82,13 +78,6 @@ public class NodeWorkspaceBuilder {
         return workspace;
     }
 
-    /**
-     * 默认构建方式:非严格隔离(向后兼容)。
-     */
-    public NodeWorkspace build(String nodeId, DagResolver.ResolvedDag dag, WorkflowContext context) {
-        return build(nodeId, dag, context, false);
-    }
-
     /**
      * 按 DAG 层级顺序排列节点 ID(先出现的层级在前)。
      */

+ 19 - 28
backend/src/main/java/com/agent/management/engine/WorkflowContext.java

@@ -8,18 +8,19 @@ import java.util.concurrent.ConcurrentHashMap;
 
 /**
  * 工作流执行上下文
- * 携带节点间传递的变量、工作目录和累积节点输出
+ *
+ * <p>变量作用域采用节点命名空间隔离:每个节点的输出存于 {@code nodeScopedOutputs[nodeId]},
+ * 同名变量互不覆盖。后继节点通过 {@link NodeInputResolver} 的显式 mapping 或四级隐式回退
+ * 从前驱命名空间取值,参见 docs/workflow-variable-scope-design.md。</p>
  */
 @Slf4j
 public class WorkflowContext {
 
     private final Long workflowId;
     private final String runId;
-    /** 扁平变量表,所有节点的输出变量合并于此(向后兼容 {{varName}} 模板渲染) */
-    private final Map<String, Object> variables = new ConcurrentHashMap<>();
     /**
      * 节点级隔离输出表:nodeId -> {varName -> value}。
-     * 用于实现变量作用域隔离,避免不同分支的同名输出相互覆盖。
+     * 同名变量按 nodeId 隔离,互不覆盖。
      */
     private final Map<String, Map<String, Object>> nodeScopedOutputs = new ConcurrentHashMap<>();
     /** 初始运行输入,与节点输出隔离保存 */
@@ -44,14 +45,10 @@ public class WorkflowContext {
         this.runId = runId;
         this.workingDir = workingDir;
         this.initialInputs = initialInputs == null ? Map.of() : new LinkedHashMap<>(initialInputs);
-        if (initialInputs != null) {
-            this.variables.putAll(initialInputs);
-        }
     }
 
     public Long getWorkflowId() { return workflowId; }
     public String getRunId() { return runId; }
-    public Map<String, Object> getVariables() { return variables; }
     public Path getWorkingDir() { return workingDir; }
     public List<NodeOutput> getNodeOutputs() { return Collections.unmodifiableList(nodeOutputs); }
     public Map<String, Object> getSections() { return sections; }
@@ -67,28 +64,13 @@ public class WorkflowContext {
         return Collections.unmodifiableMap(initialInputs);
     }
 
-    public void setVariable(String key, Object value) {
-        variables.put(key, value);
-    }
-
-    public Object getVariable(String key) {
-        return variables.get(key);
-    }
-
     /**
-     * 记录节点输出:同时更新扁平变量、节点级隔离输出和累积输出列表。
-     * 变量覆盖检测:若 key 已存在(被上游节点写过),记录 WARN 日志,便于 debug
+     * 记录节点输出:写入节点级隔离输出表与累积输出列表。
+     * 不再维护扁平变量表,同名输出按 nodeId 隔离,互不覆盖。
      */
     public void setNodeOutput(String nodeId, Map<String, Object> output) {
         if (output != null) {
             setNodeScopedOutput(nodeId, output);
-            for (Map.Entry<String, Object> e : output.entrySet()) {
-                Object existing = variables.put(e.getKey(), e.getValue());
-                if (existing != null) {
-                    log.warn("[WorkflowContext] 变量 {} 被节点 {} 覆盖(旧值类型={})",
-                            e.getKey(), nodeId, existing.getClass().getSimpleName());
-                }
-            }
             nodeOutputs.add(new NodeOutput(nodeId, output));
         }
     }
@@ -123,9 +105,18 @@ public class WorkflowContext {
     }
 
     /**
-     * 计算当前 variables 的全量快照(深拷贝),用于 debug 视图。
+     * 全量快照(debug 视图与运行历史回放)。
+     * 格式:初始输入(扁平形式) + 各节点命名空间输出({nodeId}__{varName} 形式)。
      */
-    public Map<String, Object> snapshotVariables() {
-        return new LinkedHashMap<>(variables);
+    public Map<String, Object> snapshotAllOutputs() {
+        Map<String, Object> snapshot = new LinkedHashMap<>(initialInputs);
+        for (Map.Entry<String, Map<String, Object>> nodeEntry : nodeScopedOutputs.entrySet()) {
+            String nodeId = nodeEntry.getKey();
+            for (Map.Entry<String, Object> varEntry : nodeEntry.getValue().entrySet()) {
+                snapshot.put(nodeId + NodeWorkspaceBuilder.NODE_FIELD_SEPARATOR + varEntry.getKey(),
+                        varEntry.getValue());
+            }
+        }
+        return snapshot;
     }
 }

+ 3 - 6
backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java

@@ -37,9 +37,6 @@ public class WorkflowLevelExecutor {
     private static final ObjectMapper MAPPER = new ObjectMapper()
             .registerModule(new JavaTimeModule());
 
-    /** 当前实现默认关闭严格作用域隔离,保留扁平变量以向后兼容 */
-    private static final boolean DEFAULT_SCOPE_MODE = false;
-
     private final Map<String, NodeExecutor> executorMap;
     private final SseEventBus eventBus;
     private final NodeWorkspaceBuilder workspaceBuilder;
@@ -228,7 +225,7 @@ public class WorkflowLevelExecutor {
 
         try {
             // 构建节点工作区(基于当前活跃边过滤可达前驱)
-            NodeWorkspace workspace = workspaceBuilder.build(nodeId, dag, context, activeEdges, DEFAULT_SCOPE_MODE);
+            NodeWorkspace workspace = workspaceBuilder.build(nodeId, dag, context, activeEdges);
 
             // 解析节点输入(显式 mapping + 隐式 fallback)
             NodeInputResolver.ResolvedInputs resolvedInputs = NodeInputResolver.resolveInputs(nodeData, workspace);
@@ -255,8 +252,8 @@ public class WorkflowLevelExecutor {
             if (result.getOutput() != null) {
                 context.setNodeOutput(nodeId, result.getOutput());
             }
-            // 节点完成后采集完整 variables 快照(debug 视图与运行历史回放)
-            return result.withContextSnapshot(context.snapshotVariables());
+            // 节点完成后采集完整输出快照(初始输入 + 各节点 {nodeId}__{varName} 命名空间形式,用于 debug 视图与运行历史回放)
+            return result.withContextSnapshot(context.snapshotAllOutputs());
         } catch (Exception e) {
             log.error("[WorkflowEngine] 节点 {} 执行抛出异常: {}", nodeId, e.getMessage(), e);
             return NodeExecutionResult.failed(nodeId, "节点执行异常: " + e.getMessage());

+ 2 - 2
backend/src/main/java/com/agent/management/engine/executor/OutputExecutor.java

@@ -38,8 +38,8 @@ public class OutputExecutor implements NodeExecutor {
     public NodeExecutionResult execute(String nodeId, JsonNode data, WorkflowContext context, NodeWorkspace workspace) {
         String outputType = data.path("outputType").asText("text");
 
-        // 收集上下文中的所有变量作为最终输出(保持与原 flat variables 行为一致
-        Map<String, Object> output = new HashMap<>(context.getVariables());
+        // 收集上下文所有节点输出作为最终结果(初始输入扁平 + 各节点 {nodeId}__{varName} 命名空间形式
+        Map<String, Object> output = new HashMap<>(context.snapshotAllOutputs());
 
         // 收集工作目录中的文件列表(应用 .agentignore 过滤)
         Path workingDir = context.getWorkingDir();

+ 0 - 1
backend/src/test/java/com/agent/management/engine/NodeInputResolverTest.java

@@ -18,7 +18,6 @@ class NodeInputResolverTest {
             ws.getScopedOutputs().put(e.getKey(), e.getValue());
             for (Map.Entry<String, Object> f : e.getValue().entrySet()) {
                 ws.getVariables().put(e.getKey() + NodeWorkspaceBuilder.NODE_FIELD_SEPARATOR + f.getKey(), f.getValue());
-                ws.getVariables().put(f.getKey(), f.getValue());
             }
         }
         return ws;