diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImpl.java index 1d4429f..39bf616 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImpl.java @@ -342,6 +342,13 @@ public class EdgeSystemServiceImpl implements EdgeSystemService { .setHeader(EdgeCommonUtil.buildRequest("")) .build(); SystemCommand.StopAllCommand.Feedback feedback = stub.stopAll(build); + if (feedback == null || !feedback.hasHeader()) { + throw new GlobalException("停止机器人全部设备失败: 响应为空"); + } + if (!feedback.getHeader().getSuccess()) { + throw new GlobalException("停止机器人全部设备失败: " + + feedback.getHeader().getErrorMessage()); + } return JSON.toJSONString(feedback.getHeader()); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContext.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContext.java index dc16596..600254d 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContext.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContext.java @@ -84,6 +84,9 @@ public class TaskContext { */ private volatile boolean paused = false; + /** 暂停时的设备停止和旧节点退出是否已经完成。 */ + private volatile boolean pauseSettled = true; + /** * 是否已被强行终止 */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java index 2669d4d..f09533d 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java @@ -96,10 +96,14 @@ public class TaskThreadRegistry { /** * 移除节点上下文(节点执行成功时调用) */ - public void unregisterNodeContext(String instId, String itemId, String nodeId) { + public void unregisterNodeContext(String instId, String itemId, String nodeId, + TaskNodeExecuteContext expectedContext) { String key = TaskKeyBuilder.buildKey(instId, itemId, nodeId); - nodeContextMap.remove(key); - log.info("移除节点上下文: key={}", key); + if (nodeContextMap.remove(key, expectedContext)) { + log.info("移除节点上下文: key={}", key); + } else { + log.debug("忽略过期节点上下文注销: key={}", key); + } } /** 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 8e1dac0..0e1788e 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 @@ -24,6 +24,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Collections; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; /** @@ -35,6 +36,8 @@ import java.util.stream.Collectors; @RequiredArgsConstructor public class FlowControlService { + private static final long PAUSE_SETTLE_TIMEOUT_SECONDS = 15L; + @Resource(name = "threadPoolTaskExecutor") private final ThreadPoolTaskExecutor executor; private final TaskInstHolder taskInstHolder; @@ -80,15 +83,22 @@ public class FlowControlService { if (ctx.isPaused()) { throw new GlobalException("任务已处于暂停状态"); } - // 异步停止终端 - stopAllRobots(ctx); ctx.setPaused(true); + ctx.setPauseSettled(false); ctx.setTerminalStatus(TerminalStatusEnum.RUNNING); ctx.setStatus(TaskStatusEnum.PAUSED); ctx.updateTimestamp(); + List interruptedNodes = taskThreadRegistry.getPendingNodes(instId); + taskThreadRegistry.interruptAll(instId); + // StopAll 是阻塞式调用。只有设备侧停止完成,暂停接口才返回,避免恢复请求和旧动作并发。 + boolean robotsStopped = stopAllRobotsAndWait(ctx); + boolean nodesFinished = waitForInterruptedNodes(instId, interruptedNodes); + ctx.setPauseSettled(robotsStopped && nodesFinished); + if (!ctx.isPauseSettled()) { + log.warn("任务已进入暂停状态,但设备停止仍在收敛: instId={}", instId); + } // 统一记录日志 + 设置上下文状态 + 数据库状态 taskInstHolder.syncStatus(instId, TaskStatusEnum.PAUSED); - taskThreadRegistry.interruptAll(instId); // 被中断的节点 状态修改成 PAUSED for (TaskNodeExecuteContext pendingNode : taskThreadRegistry.getPendingNodes(instId)) { String nodeId = pendingNode.getNode().getNodeId(); @@ -107,6 +117,15 @@ public class FlowControlService { if (!ctx.isPaused()) { throw new GlobalException("当前任务未处于暂停状态,无法继续运行"); } + if (!ctx.isPauseSettled()) { + List interruptedNodes = taskThreadRegistry.getPendingNodes(instId); + boolean robotsStopped = stopAllRobotsAndWait(ctx); + boolean nodesFinished = waitForInterruptedNodes(instId, interruptedNodes); + ctx.setPauseSettled(robotsStopped && nodesFinished); + if (!ctx.isPauseSettled()) { + throw new GlobalException("设备停止尚未完成,请稍后重新恢复"); + } + } ctx.setPaused(false); @@ -207,4 +226,48 @@ public class FlowControlService { .distinct() .forEach(robotId -> executor.execute(() -> edgeSystemService.stopAll(robotId))); } + + private boolean stopAllRobotsAndWait(TaskContext context) { + List robotIds = context.getRobotIds() == null || context.getRobotIds().isEmpty() + ? Collections.singletonList(context.getRobotId()) + : context.getRobotIds(); + boolean success = true; + for (String robotId : robotIds.stream() + .filter(robotId -> robotId != null && !robotId.isBlank()) + .distinct() + .collect(Collectors.toList())) { + try { + edgeSystemService.stopAll(robotId); + } catch (RuntimeException exception) { + success = false; + log.warn("暂停机器人设备失败,恢复时将重试: robotId={}, error={}", + robotId, exception.getMessage()); + } + } + return success; + } + + private boolean waitForInterruptedNodes(String instId, List interruptedNodes) { + long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(PAUSE_SETTLE_TIMEOUT_SECONDS); + for (TaskNodeExecuteContext pendingNode : interruptedNodes) { + if (pendingNode.getNode().getNodeType() == NodeTypeEnum.LOOP) { + continue; + } + long remainingNanos = deadlineNanos - System.nanoTime(); + if (remainingNanos <= 0) { + log.warn("暂停等待节点退出超时: instId={}, nodeId={}", instId, pendingNode.getNode().getNodeId()); + return false; + } + try { + if (!pendingNode.awaitExecutionFinished(remainingNanos, TimeUnit.NANOSECONDS)) { + log.warn("暂停等待节点退出超时: instId={}, nodeId={}", instId, pendingNode.getNode().getNodeId()); + return false; + } + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new GlobalException("暂停任务时等待设备动作结束被中断"); + } + } + return true; + } } 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 249621f..0729953 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 @@ -118,7 +118,9 @@ public class FlowItemExecutor { log.info("当前执行的是 {} 节点,迭代路径={}", node.getNodeName(), iterations); //注册当前线程 - registerThread(graph, node, startNodeId, startInputDefs, rootMessage, onFinished, iterations, instId, itemId, nodeId); + TaskNodeExecuteContext nodeExecutionContext = registerThread( + graph, node, startNodeId, startInputDefs, rootMessage, onFinished, + iterations, instId, itemId, nodeId); try { // 创建带上下文的执行消息 @@ -172,10 +174,11 @@ public class FlowItemExecutor { } throw e; } finally { + nodeExecutionContext.markExecutionFinished(); // 如果不是暂停全中断 TaskContext ctx = taskInstHolder.getContext(instId); if (ctx == null || !ctx.isPaused()) { - taskThreadRegistry.unregisterNodeContext(instId, itemId, nodeId); + taskThreadRegistry.unregisterNodeContext(instId, itemId, nodeId, nodeExecutionContext); } } } @@ -286,7 +289,7 @@ public class FlowItemExecutor { ); } - private void registerThread(FlowGraph graph, FlowNodeWrapper node, String startNodeId, List startInputDefs, + private TaskNodeExecuteContext registerThread(FlowGraph graph, FlowNodeWrapper node, String startNodeId, List startInputDefs, TaskNodeExecuteMessage rootMessage, Runnable onFinished, List iterations, String instId, String itemId, String nodeId) { TaskNodeExecuteContext nodeExecuteContext = new TaskNodeExecuteContext(); nodeExecuteContext.setThread(Thread.currentThread()); @@ -299,6 +302,7 @@ public class FlowItemExecutor { // nodeExecuteContext.setLoopIteration(loopIteration); nodeExecuteContext.setIterations(new ArrayList<>(iterations)); taskThreadRegistry.registerNodeContext(instId, itemId, nodeId, nodeExecuteContext); + return nodeExecuteContext; } /** 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 e60597a..c60127a 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 @@ -7,12 +7,16 @@ import lombok.Data; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; /** * 节点执行上下文,用于保存线程与流程信息 */ @Data public class TaskNodeExecuteContext { + private final CountDownLatch executionFinished = new CountDownLatch(1); + private Thread thread; private FlowGraph graph; @@ -24,4 +28,12 @@ public class TaskNodeExecuteContext { private int loopIteration; private Object loopArray; private List iterations = new ArrayList<>(); + + public void markExecutionFinished() { + executionFinished.countDown(); + } + + public boolean awaitExecutionFinished(long timeout, TimeUnit unit) throws InterruptedException { + return executionFinished.await(timeout, unit); + } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java new file mode 100644 index 0000000..f43ac28 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java @@ -0,0 +1,24 @@ +package com.cmvr.test.flow.context; + +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteContext; +import org.junit.Test; + +import static org.junit.Assert.assertSame; + +public class TaskThreadRegistryTest { + + @Test + public void staleExecutionCannotUnregisterResumedNodeContext() { + TaskThreadRegistry registry = new TaskThreadRegistry(); + TaskNodeExecuteContext interrupted = new TaskNodeExecuteContext(); + TaskNodeExecuteContext resumed = new TaskNodeExecuteContext(); + interrupted.setThread(Thread.currentThread()); + resumed.setThread(Thread.currentThread()); + + registry.registerNodeContext("inst-1", "item-1", "node-1", interrupted); + registry.registerNodeContext("inst-1", "item-1", "node-1", resumed); + registry.unregisterNodeContext("inst-1", "item-1", "node-1", interrupted); + + assertSame(resumed, registry.getNodeContext("inst-1", "item-1", "node-1")); + } +}