fix(flow): 解决任务暂停时设备停止和节点退出的同步问题

- 在EdgeSystemServiceImpl中添加响应验证和异常处理
- 实现暂停操作的阻塞等待机制确保设备完全停止
- 添加节点执行完成标记和等待功能
- 优化任务上下文中的暂停状态管理
- 防止恢复操作时新旧节点上下文冲突
- 添加超时控制和异常重试机制
This commit is contained in:
lixiaolong 2026-08-13 10:12:14 +08:00
parent 05992d1980
commit 6265db01c7
7 changed files with 126 additions and 9 deletions

View File

@ -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());
}

View File

@ -84,6 +84,9 @@ public class TaskContext {
*/
private volatile boolean paused = false;
/** 暂停时的设备停止和旧节点退出是否已经完成。 */
private volatile boolean pauseSettled = true;
/**
* 是否已被强行终止
*/

View File

@ -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);
if (nodeContextMap.remove(key, expectedContext)) {
log.info("移除节点上下文: key={}", key);
} else {
log.debug("忽略过期节点上下文注销: key={}", key);
}
}
/**

View File

@ -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<TaskNodeExecuteContext> 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<TaskNodeExecuteContext> 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<String> 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<TaskNodeExecuteContext> 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;
}
}

View File

@ -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<FlowParamDef> startInputDefs,
private TaskNodeExecuteContext registerThread(FlowGraph graph, FlowNodeWrapper node, String startNodeId, List<FlowParamDef> startInputDefs,
TaskNodeExecuteMessage rootMessage, Runnable onFinished, List<Integer> 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;
}
/**

View File

@ -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<Integer> iterations = new ArrayList<>();
public void markExecutionFinished() {
executionFinished.countDown();
}
public boolean awaitExecutionFinished(long timeout, TimeUnit unit) throws InterruptedException {
return executionFinished.await(timeout, unit);
}
}

View File

@ -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"));
}
}