From e40bd9dc82811cea09eadfbdc28da2c68104221d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=99=88=E5=AD=90=E9=BB=98?= <925456043@qq.com> Date: Mon, 31 Aug 2026 15:54:50 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20M28=20=E5=A2=9E=E5=8A=A0=E5=B7=A5?= =?UTF-8?q?=E4=BD=9C=E6=B5=81=E6=B1=87=E8=81=9A=E5=AE=89=E5=85=A8=E6=A0=A1?= =?UTF-8?q?=E9=AA=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/WorkflowCheckService.java | 163 +++++++++++++++++ .../service/WorkflowCheckServiceTest.java | 166 ++++++++++++++++++ 2 files changed, 329 insertions(+) diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckService.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckService.java index 96489d23..8211a68f 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckService.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckService.java @@ -59,6 +59,8 @@ public class WorkflowCheckService { private static final String SYSTEM_START_PARAM_NAME = "user_input"; private static final int MIN_LOOP_COUNT = 1; private static final int MAX_LOOP_COUNT = 300; + private static final String JOIN_MODE_ANY = "any"; + private static final String JOIN_MODE_ALL = "all"; @Resource private WorkflowService workflowService; @@ -196,6 +198,10 @@ public class WorkflowCheckService { edge.id = trimToNull(edgeJson.getString("id")); edge.source = trimToNull(edgeJson.getString("source")); edge.target = trimToNull(edgeJson.getString("target")); + JSONObject edgeData = edgeJson.getJSONObject("data"); + edge.condition = edgeData == null + ? null + : trimToNull(edgeData.getString("condition")); if (!StringUtils.hasText(edge.id)) { addIssue(issues, issueKeys, "EDGE_ID_EMPTY", "存在连线缺少 id", null, null, null); @@ -228,10 +234,162 @@ public class WorkflowCheckService { parsedWorkflow.nodes = nodes; parsedWorkflow.edges = edges; parsedWorkflow.nodeMap = nodeMap; + checkJoinModes(parsedWorkflow, issues, issueKeys); checkDatacenterNodes(parsedWorkflow, issues, issueKeys); return parsedWorkflow; } + /** + * 校验节点汇聚模式及其静态可证明的到达安全性。 + * + * @param parsed 工作流视图 + * @param issues 问题列表 + * @param issueKeys 问题去重键 + */ + private void checkJoinModes( + ParsedWorkflow parsed, + List issues, + Set issueKeys) { + Map> inwardEdges = new LinkedHashMap<>(); + for (EdgeView edge : parsed.edges) { + if (edge == null || !StringUtils.hasText(edge.target)) { + continue; + } + inwardEdges.computeIfAbsent( + edge.target, ignored -> new ArrayList<>()).add(edge); + } + + for (NodeView node : parsed.nodes) { + String joinMode = resolveJoinMode(node); + if (joinMode == null) { + addIssue( + issues, + issueKeys, + "JOIN_MODE_INVALID", + "执行时机配置无效,joinMode 仅支持 any 或 all", + node.id, + null, + node.name); + continue; + } + if (JOIN_MODE_ALL.equals(joinMode) + && StringUtils.hasText(node.parentId)) { + addIssue( + issues, + issueKeys, + "JOIN_MODE_LOOP_CHILD_UNSUPPORTED", + "显式循环子图暂不支持“全部上游完成”,请改为“任一上游完成”", + node.id, + null, + node.name); + } + } + + Set guaranteedNodes = findGuaranteedNodes( + parsed, inwardEdges); + for (NodeView node : parsed.nodes) { + if (!JOIN_MODE_ALL.equals(resolveJoinMode(node)) + || StringUtils.hasText(node.parentId)) { + continue; + } + List directInward = inwardEdges.getOrDefault( + node.id, Collections.emptyList()); + if (directInward.size() <= 1) { + continue; + } + boolean allGuaranteed = directInward.stream().allMatch(edge -> + !edge.hasCondition() + && guaranteedNodes.contains(edge.source)); + if (!allGuaranteed) { + addIssue( + issues, + issueKeys, + "JOIN_MODE_CONDITIONAL_PATH_UNSUPPORTED", + "“全部上游完成”可能永久等待:存在条件、互斥或无法证明必达的上游路径。" + + "请改为“任一上游完成”或调整连线,确保所有直接入边都会到达", + node.id, + null, + node.name); + } + } + } + + /** + * 使用保守固定点传播计算能够保证执行的根级节点。 + * + * @param parsed 工作流视图 + * @param inwardEdges 直接入边索引 + * @return 保证执行的节点 ID + */ + private Set findGuaranteedNodes( + ParsedWorkflow parsed, + Map> inwardEdges) { + Set guaranteed = parsed.nodes.stream() + .filter(NodeView::isRootLevel) + .filter(node -> TYPE_START.equals(node.type)) + .map(node -> node.id) + .filter(StringUtils::hasText) + .collect(Collectors.toCollection(LinkedHashSet::new)); + + boolean changed; + do { + changed = false; + for (NodeView node : parsed.nodes) { + if (!node.isRootLevel() + || guaranteed.contains(node.id) + || hasAdvancedCondition(node)) { + continue; + } + String joinMode = resolveJoinMode(node); + if (joinMode == null) { + continue; + } + List directInward = inwardEdges.getOrDefault( + node.id, Collections.emptyList()); + boolean isGuaranteed; + if (JOIN_MODE_ALL.equals(joinMode)) { + isGuaranteed = !directInward.isEmpty() + && directInward.stream().allMatch(edge -> + !edge.hasCondition() + && guaranteed.contains(edge.source)); + } else { + isGuaranteed = directInward.stream().anyMatch(edge -> + !edge.hasCondition() + && guaranteed.contains(edge.source)); + } + if (isGuaranteed && guaranteed.add(node.id)) { + changed = true; + } + } + } while (changed); + return guaranteed; + } + + /** + * 读取节点汇聚模式。字段缺失时兼容为 any,显式非法值返回 null。 + */ + private String resolveJoinMode(NodeView node) { + if (node == null || node.data == null + || !node.data.containsKey("joinMode")) { + return JOIN_MODE_ANY; + } + String value = trimToNull(node.data.getString("joinMode")); + if (JOIN_MODE_ANY.equalsIgnoreCase(value)) { + return JOIN_MODE_ANY; + } + if (JOIN_MODE_ALL.equalsIgnoreCase(value)) { + return JOIN_MODE_ALL; + } + return null; + } + + private boolean hasAdvancedCondition(NodeView node) { + return node != null + && node.data != null + && StringUtils.hasText( + trimToNull(node.data.getString("condition"))); + } + /** * 校验普通循环、显式循环和循环父子层级。 * @@ -1610,5 +1768,10 @@ public class WorkflowCheckService { private String id; private String source; private String target; + private String condition; + + private boolean hasCondition() { + return StringUtils.hasText(condition); + } } } diff --git a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckServiceTest.java b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckServiceTest.java index 2be7f0ff..e903c5f0 100644 --- a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckServiceTest.java +++ b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/service/WorkflowCheckServiceTest.java @@ -22,6 +22,163 @@ import java.util.Map; public class WorkflowCheckServiceTest { + @Test + public void testSaveAndPreExecuteShouldPassGuaranteedAllJoin() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject joinData = data("汇聚"); + joinData.put("joinMode", "all"); + String content = workflowJson( + array( + node("start", "startNode", null, data("开始")), + node("a", "codeNode", null, data("分支 A")), + node("b", "codeNode", null, data("分支 B")), + node("join", "codeNode", null, joinData), + node("end", "endNode", null, data("结束"))), + array( + edge("start-a", "start", "a"), + edge("start-b", "start", "b"), + edge("a-join", "a", "join"), + edge("b-join", "b", "join"), + edge("join-end", "join", "end"))); + + Assert.assertTrue(service.checkContent( + content, WorkflowCheckStage.SAVE, null).isPassed()); + Assert.assertTrue(service.checkContent( + content, WorkflowCheckStage.PRE_EXECUTE, null).isPassed()); + } + + @Test + public void testSaveAndPreExecuteShouldBlockConditionalAllJoin() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject joinData = data("汇聚"); + joinData.put("joinMode", "all"); + String content = workflowJson( + array( + node("start", "startNode", null, data("开始")), + node("a", "codeNode", null, data("条件来源")), + node("b", "codeNode", null, data("普通来源")), + node("join", "codeNode", null, joinData), + node("end", "endNode", null, data("结束"))), + array( + conditionalEdge("start-a", "start", "a", "enabled === true"), + edge("start-b", "start", "b"), + edge("a-join", "a", "join"), + edge("b-join", "b", "join"), + edge("join-end", "join", "end"))); + + WorkflowCheckResult save = service.checkContent( + content, WorkflowCheckStage.SAVE, null); + WorkflowCheckResult preExecute = service.checkContent( + content, WorkflowCheckStage.PRE_EXECUTE, null); + + Assert.assertFalse(save.isPassed()); + Assert.assertFalse(preExecute.isPassed()); + assertHasCode(save, "JOIN_MODE_CONDITIONAL_PATH_UNSUPPORTED"); + assertHasCode(preExecute, "JOIN_MODE_CONDITIONAL_PATH_UNSUPPORTED"); + Assert.assertTrue(save.getIssues().stream().anyMatch(issue -> + "join".equals(issue.getNodeId()) + && issue.getMessage().contains("永久等待") + && issue.getMessage().contains("任一上游完成"))); + } + + @Test + public void testSaveShouldBlockAllJoinWithDirectConditionalEdge() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject joinData = data("汇聚"); + joinData.put("joinMode", "all"); + String content = workflowJson( + array( + node("start", "startNode", null, data("开始")), + node("a", "codeNode", null, data("分支 A")), + node("b", "codeNode", null, data("分支 B")), + node("join", "codeNode", null, joinData)), + array( + edge("start-a", "start", "a"), + edge("start-b", "start", "b"), + conditionalEdge("a-join", "a", "join", "matched === true"), + edge("b-join", "b", "join"))); + + WorkflowCheckResult result = service.checkContent( + content, WorkflowCheckStage.SAVE, null); + + Assert.assertFalse(result.isPassed()); + assertHasCode(result, "JOIN_MODE_CONDITIONAL_PATH_UNSUPPORTED"); + } + + @Test + public void testSaveShouldBlockAllJoinFromCustomConditionSource() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject conditionalSource = data("高级条件来源"); + conditionalSource.put("condition", "score > 0"); + JSONObject joinData = data("汇聚"); + joinData.put("joinMode", "all"); + String content = workflowJson( + array( + node("start", "startNode", null, data("开始")), + node("a", "codeNode", null, conditionalSource), + node("b", "codeNode", null, data("普通来源")), + node("join", "codeNode", null, joinData)), + array( + edge("start-a", "start", "a"), + edge("start-b", "start", "b"), + edge("a-join", "a", "join"), + edge("b-join", "b", "join"))); + + WorkflowCheckResult result = service.checkContent( + content, WorkflowCheckStage.SAVE, null); + + Assert.assertFalse(result.isPassed()); + assertHasCode(result, "JOIN_MODE_CONDITIONAL_PATH_UNSUPPORTED"); + } + + @Test + public void testSaveShouldBlockInvalidAndLoopChildJoinModes() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject invalidData = data("非法汇聚"); + invalidData.put("joinMode", "first"); + JSONObject loopData = loopData( + fixedParameter("count", "2", "Number"), null); + JSONObject childData = data("循环子节点"); + childData.put("joinMode", "all"); + String content = workflowJson( + array( + node("invalid", "codeNode", null, invalidData), + node("loop", "loopNode", null, loopData), + node("child", "codeNode", "loop", childData)), + new JSONArray()); + + WorkflowCheckResult result = service.checkContent( + content, WorkflowCheckStage.SAVE, null); + + Assert.assertFalse(result.isPassed()); + assertHasCode(result, "JOIN_MODE_INVALID"); + assertHasCode(result, "JOIN_MODE_LOOP_CHILD_UNSUPPORTED"); + Assert.assertTrue(result.getIssues().stream().anyMatch(issue -> + "invalid".equals(issue.getNodeId()) + && "JOIN_MODE_INVALID".equals(issue.getCode()))); + Assert.assertTrue(result.getIssues().stream().anyMatch(issue -> + "child".equals(issue.getNodeId()) + && "JOIN_MODE_LOOP_CHILD_UNSUPPORTED".equals(issue.getCode()))); + } + + @Test + public void testSaveShouldAllowSingleConditionalInboundAllJoin() throws Exception { + WorkflowCheckService service = newService(new HashMap<>()); + JSONObject joinData = data("单入边汇聚"); + joinData.put("joinMode", "all"); + String content = workflowJson( + array( + node("start", "startNode", null, data("开始")), + node("join", "codeNode", null, joinData)), + array(conditionalEdge( + "start-join", "start", "join", "enabled === true"))); + + WorkflowCheckResult result = service.checkContent( + content, WorkflowCheckStage.SAVE, null); + + Assert.assertTrue(result.isPassed()); + } + /** * 验证保存阶段接受合法的正则条件规则。 */ @@ -992,4 +1149,13 @@ public class WorkflowCheckServiceTest { edge.put("target", target); return edge; } + + private static JSONObject conditionalEdge( + String id, String source, String target, String condition) { + JSONObject edge = edge(id, source, target); + JSONObject data = new JSONObject(); + data.put("condition", condition); + edge.put("data", data); + return edge; + } }