From a1c30402104f88ce76d126ee5d6787414f0b90e5 Mon Sep 17 00:00:00 2001 From: stream Date: Tue, 24 Mar 2026 10:59:50 +0800 Subject: [PATCH] =?UTF-8?q?fix(flow):=20=E4=BF=AE=E5=A4=8D=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E6=8E=A7=E5=88=B6=E4=B8=AD=E7=9A=84=E5=BE=AA=E7=8E=AF?= =?UTF-8?q?=E6=95=B0=E7=BB=84=E4=BC=A0=E9=80=92=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 在EdgeHlcServiceImpl中添加触控坐标验证,防止坐标为0时的异常 - 在FlowControlService中为resumeMessage添加loopArray属性传递,确保恢复时保留循环数组数据 - 修改FlowGetCurrentObjNodeHandler中获取当前对象的方式,使用loopArray替代inputParams中的array - 在FlowItemExecutor中将rootMessage的loopArray传递给子消息 - 更新FlowLoopNodeHandler中循环处理逻辑,确保loopArray在循环过程中正确传递 - 调整FlowNodeParamPreparer中循环目标解析方式,直接使用loopNumVal作为目标对象 - 在TaskNodeExecuteContext和TaskNodeExecuteMessage中添加loopArray属性支持 --- .../service/impl/EdgeHlcServiceImpl.java | 4 ++++ .../test/flow/control/FlowControlService.java | 8 ++++++- .../FlowGetCurrentObjNodeHandler.java | 21 ++++++++++--------- .../dispatcher/FlowLoopNodeHandler.java | 4 +++- .../flow/runtime/engine/FlowItemExecutor.java | 4 +++- .../engine/support/FlowNodeParamPreparer.java | 5 +++-- .../message/TaskNodeExecuteContext.java | 1 + .../message/TaskNodeExecuteMessage.java | 2 ++ .../impl/TiVehicleFunctionServiceImpl.java | 2 +- 9 files changed, 35 insertions(+), 16 deletions(-) diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java index 4b28786..ff9dbe7 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java @@ -3,6 +3,7 @@ package com.cmvr.edge.client.service.impl; import cmvr.api.HlcCommand; import cmvr.api.HlcServiceGrpc; import com.alibaba.fastjson2.JSON; +import com.cmvr.common.exception.GlobalException; import com.cmvr.edge.client.manage.GrpcServiceManager; import com.cmvr.edge.client.model.hlc.EdgeTouchVO; import com.cmvr.edge.client.service.EdgeHlcService; @@ -85,6 +86,9 @@ public class EdgeHlcServiceImpl implements EdgeHlcService { @Override public String touch(EdgeTouchVO edgeTouchVO) { + if (edgeTouchVO.getX() == edgeTouchVO.getY() && edgeTouchVO.getX() == 0) { + throw new GlobalException("触控坐标不能为0"); + } HlcServiceGrpc.HlcServiceBlockingStub stub = grpcServiceManager.getGrpcClient(edgeTouchVO.getTerminalId(), HlcServiceGrpc.HlcServiceBlockingStub.class); HlcCommand.Touch.Request request = HlcCommand.Touch.Request.newBuilder() .setHeader(EdgeCommonUtil.buildRequest(edgeTouchVO.getDeviceId())) diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java index c9af14c..19c6bd3 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java @@ -78,7 +78,11 @@ public class FlowControlService { if (ctx.isPaused()) { throw new GlobalException("任务已处于暂停状态"); } - edgeSystemService.stopAll(ctx.getTerminalId()); + try { + edgeSystemService.stopAll(ctx.getTerminalId()); + } catch (Exception e) { + + } ctx.setPaused(true); ctx.setTerminalStatus(TerminalStatusEnum.RUNNING); ctx.setStatus(TaskStatusEnum.PAUSED); @@ -144,6 +148,7 @@ public class FlowControlService { TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage); resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留中断时的循环次数 + resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留中断时的循环次数 resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留中断时的路径 flowItemExecutor.executeNode( @@ -169,6 +174,7 @@ public class FlowControlService { TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage); resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留循环次数 + resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留循环次数 resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留路径 flowItemExecutor.executeNode( diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java index bee6c8e..09ff6f8 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java @@ -8,6 +8,7 @@ import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; +import java.util.Collections; import java.util.List; @Slf4j @@ -18,17 +19,17 @@ public class FlowGetCurrentObjNodeHandler implements FlowNodeTypeHandler { @Override public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) { try { - JSONObject inputParams = message.getInputParams(); - JSONArray jsonArray = inputParams.getJSONArray("array"); - - List iterations = message.getIterations(); - int last = CollUtil.getLast(iterations) - 1; - +// JSONObject inputParams = message.getInputParams(); +// JSONArray jsonArray = inputParams.getJSONArray("array"); +// +// List iterations = message.getIterations(); +// int last = CollUtil.getLast(iterations) - 1; + int index = message.getLoopNum() - 1; JSONObject output = new JSONObject(); - output.put("index", last); - - if (CollUtil.isNotEmpty(jsonArray) && last >= 0) { - output.put("object", jsonArray.get(last)); + output.put("index", index); + JSONArray objects = (JSONArray)message.getLoopArray(); + if (CollUtil.isNotEmpty(objects) && index >= 0) { + output.put("object", objects.get(index)); } return TaskNodeExecuteResult.success(output); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java index 7b11a01..a9e5bb5 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java @@ -37,6 +37,7 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler { String nodeId = message.getNodeId(); FlowNodeWrapper nodeWrapper = graph.getNode(nodeId); int loopCount = inputParams.getIntValue("loopNum"); // 获取循环次数 + Object loopArray = inputParams.get("loopArray"); // 获取循环次数 // 如果没有设置循环次数或者循环次数为 0,则跳过处理 if (loopCount <= 0) { @@ -65,8 +66,9 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler { // 克隆 message,并明确设置 loopNum TaskNodeExecuteMessage subMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(message, subMessage); -// subMessage.setLoopNum(i); // 当前 loop 第 i 次 + subMessage.setLoopNum(i); // 当前 loop 第 i 次 subMessage.setIterations(newIterations); // 完整路径 + subMessage.setLoopArray(loopArray); flowItemExecutor.executeSubGraph(subGraph, subMessage, latch::countDown, newIterations); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java index 1eac761..02dbf09 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java @@ -84,7 +84,7 @@ public class FlowItemExecutor { // 执行子图起始节点 executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs, rootMessage, onFinished, iterations); - + log.info("当前执行-子图"); // 启动子图调度 flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()), executor, onFinished, iterations); @@ -172,6 +172,8 @@ public class FlowItemExecutor { // loopNum 不在这里计算,而是由 FlowLoopNodeHandler / resume 显式写入 // 如果 rootMessage 里已经带了 loopNum,就沿用它 message.setLoopNum(rootMessage.getLoopNum()); + + message.setLoopArray(rootMessage.getLoopArray()); return message; } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java index 78a49b7..cebde64 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java @@ -84,11 +84,12 @@ public class FlowNodeParamPreparer { // 4 LOOP 节点:动态计算 loopCount if (node.getNodeType() == NodeTypeEnum.LOOP) { Object loopNumVal = input.get("loopNum"); + input.put("loopArray", loopNumVal); int loopCount = 0; // 根据迭代路径找到当前层的集合对象 - Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0); - +// Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0); + Object target = loopNumVal; if (target instanceof JSONArray) { loopCount = ((JSONArray) target).size(); } else if (target instanceof Collection) { diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java index 5c2f57c..e60597a 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java @@ -22,5 +22,6 @@ public class TaskNodeExecuteContext { private TaskNodeExecuteMessage rootMessage; private Runnable onFinished; private int loopIteration; + private Object loopArray; private List iterations = new ArrayList<>(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java index b887a37..33462f3 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java @@ -55,6 +55,8 @@ public class TaskNodeExecuteMessage { */ private List iterations = new ArrayList<>(); + private Object loopArray; + /** * 节点名称 */ diff --git a/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java b/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java index 5681d03..fcfb7e0 100644 --- a/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java +++ b/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java @@ -40,7 +40,7 @@ public class TiVehicleFunctionServiceImpl extends ServiceImpl