feat(flow): 添加流程失败清理服务和机械臂使能枚举

- 新增 FlowFailureCleanupService 用于清理设备端和后端媒体状态
- 添加 ARM_ENABLE 和 ARM_DISABLE 枚举值支持机械臂使能控制
- 在 TaskInstHolder 中集成失败清理逻辑
- 实现机器人录音、录像和系统状态的异常清理功能
- 添加 FlowFailureCleanupService 单元测试验证清理逻辑
- 注入 FlowFailureCleanupService 并在任务失败时执行清理操作
This commit is contained in:
lixiaolong 2026-08-18 10:19:25 +08:00
parent d3498e0178
commit 0b05438712
4 changed files with 119 additions and 0 deletions

View File

@ -65,6 +65,10 @@ public enum ActionEnum {
// 末端运动 // 末端运动
ARM_MOVE_TO_POINT("EDGE", "ARM_MOVE_TO_POINT", "运动到指定点"), ARM_MOVE_TO_POINT("EDGE", "ARM_MOVE_TO_POINT", "运动到指定点"),
ARM_MOVE_TO_J("EDGE", "ARM_MOVE_TO_J", "关节空间运动"), ARM_MOVE_TO_J("EDGE", "ARM_MOVE_TO_J", "关节空间运动"),
// 开始使能
ARM_ENABLE("EDGE", "ARM_ENABLE", "使能"),
// 停止使能
ARM_DISABLE("EDGE", "ARM_DISABLE", "停止使能"),
// --------------- 头部 --------------- // --------------- 头部 ---------------
BIO_HEAD_SPEAK_START("EDGE", "BIO_HEAD_SPEAK_START", "开始说话"), BIO_HEAD_SPEAK_START("EDGE", "BIO_HEAD_SPEAK_START", "开始说话"),

View File

@ -4,12 +4,14 @@ import com.alibaba.fastjson2.JSONObject;
import com.cmvr.test.enums.NodeTypeEnum; import com.cmvr.test.enums.NodeTypeEnum;
import com.cmvr.test.enums.TaskStatusEnum; import com.cmvr.test.enums.TaskStatusEnum;
import com.cmvr.test.enums.TerminalStatusEnum; import com.cmvr.test.enums.TerminalStatusEnum;
import com.cmvr.test.flow.control.FlowFailureCleanupService;
import com.cmvr.test.flow.runtime.event.FlowExecutionEvent; import com.cmvr.test.flow.runtime.event.FlowExecutionEvent;
import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher; import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher;
import com.cmvr.test.service.ITeNodeInstService; import com.cmvr.test.service.ITeNodeInstService;
import com.cmvr.test.service.ITeTaskInstService; import com.cmvr.test.service.ITeTaskInstService;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import jakarta.annotation.Resource;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.util.List; import java.util.List;
@ -22,6 +24,9 @@ import java.util.List;
@RequiredArgsConstructor @RequiredArgsConstructor
public class TaskInstHolder { public class TaskInstHolder {
@Resource
private FlowFailureCleanupService flowFailureCleanupService;
private final ITeNodeInstService nodeInstService; private final ITeNodeInstService nodeInstService;
private final TaskContextManager taskContextManager; private final TaskContextManager taskContextManager;
private final ITeTaskInstService taskInstService; private final ITeTaskInstService taskInstService;
@ -123,6 +128,9 @@ public class TaskInstHolder {
// 终止和并发失败都可能先注销上下文,失败状态仍需落库。 // 终止和并发失败都可能先注销上下文,失败状态仍需落库。
TaskContext ctx = taskContextManager.get(instId); TaskContext ctx = taskContextManager.get(instId);
String moduleCode = ctx != null ? ctx.getModuleCode() : null; String moduleCode = ctx != null ? ctx.getModuleCode() : null;
if (ctx != null && flowFailureCleanupService != null) {
flowFailureCleanupService.cleanup(ctx);
}
syncStatus(instId, TaskStatusEnum.FAILED); syncStatus(instId, TaskStatusEnum.FAILED);
// 注销任务上下文 // 注销任务上下文

View File

@ -0,0 +1,60 @@
package com.cmvr.test.flow.control;
import com.cmvr.edge.client.service.EdgeCameraService;
import com.cmvr.edge.client.service.EdgeMicrophoneService;
import com.cmvr.edge.client.service.EdgeSystemService;
import com.cmvr.test.flow.context.TaskContext;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
/**
* Cleans device-side and backend-local media state when a flow fails unexpectedly.
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class FlowFailureCleanupService {
private final EdgeMicrophoneService edgeMicrophoneService;
private final EdgeCameraService edgeCameraService;
private final EdgeSystemService edgeSystemService;
public void cleanup(TaskContext context) {
if (context == null) {
return;
}
context.setStopped(true);
List<String> robotIds = context.getRobotIds() == null || context.getRobotIds().isEmpty()
? Collections.singletonList(context.getRobotId())
: context.getRobotIds();
robotIds.stream()
.filter(Objects::nonNull)
.map(String::trim)
.filter(robotId -> !robotId.isEmpty())
.distinct()
.forEach(this::cleanupRobot);
}
private void cleanupRobot(String robotId) {
try {
edgeMicrophoneService.abortWorkflowRecordings(robotId);
} catch (RuntimeException exception) {
log.warn("清理失败任务录音会话失败: robotId={}, error={}", robotId, exception.getMessage());
}
try {
edgeCameraService.abortWorkflowRecordings(robotId);
} catch (RuntimeException exception) {
log.warn("清理失败任务录像会话失败: robotId={}, error={}", robotId, exception.getMessage());
}
try {
edgeSystemService.stopAll(robotId);
} catch (RuntimeException exception) {
log.warn("停止失败任务设备失败: robotId={}, error={}", robotId, exception.getMessage());
}
}
}

View File

@ -0,0 +1,47 @@
package com.cmvr.test.flow.control;
import com.cmvr.edge.client.service.EdgeCameraService;
import com.cmvr.edge.client.service.EdgeMicrophoneService;
import com.cmvr.edge.client.service.EdgeSystemService;
import com.cmvr.test.flow.context.TaskContext;
import org.junit.Test;
import java.lang.reflect.Proxy;
import java.util.ArrayList;
import java.util.List;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
public class FlowFailureCleanupServiceTest {
@Test
public void clearsMediaAndStopsEveryRobotOnFailure() {
List<String> calls = new ArrayList<>();
EdgeMicrophoneService microphone = proxy(EdgeMicrophoneService.class, "microphone", calls);
EdgeCameraService camera = proxy(EdgeCameraService.class, "camera", calls);
EdgeSystemService system = proxy(EdgeSystemService.class, "system", calls);
FlowFailureCleanupService service = new FlowFailureCleanupService(microphone, camera, system);
TaskContext context = new TaskContext();
context.setRobotIds(List.of("robot-1", "robot-2"));
service.cleanup(context);
assertTrue(context.isStopped());
assertEquals(List.of(
"microphone:robot-1", "camera:robot-1", "system:robot-1",
"microphone:robot-2", "camera:robot-2", "system:robot-2"), calls);
}
@SuppressWarnings("unchecked")
private <T> T proxy(Class<T> type, String name, List<String> calls) {
return (T) Proxy.newProxyInstance(type.getClassLoader(), new Class<?>[]{type},
(proxy, method, args) -> {
if ("abortWorkflowRecordings".equals(method.getName())
|| "stopAll".equals(method.getName())) {
calls.add(name + ":" + args[0]);
}
return null;
});
}
}