From 0b05438712e5d0d532ab591887608dd13c760e38 Mon Sep 17 00:00:00 2001 From: lixiaolong <702156524@qq.com> Date: Tue, 18 Aug 2026 10:19:25 +0800 Subject: [PATCH] =?UTF-8?q?feat(flow):=20=E6=B7=BB=E5=8A=A0=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E5=A4=B1=E8=B4=A5=E6=B8=85=E7=90=86=E6=9C=8D=E5=8A=A1?= =?UTF-8?q?=E5=92=8C=E6=9C=BA=E6=A2=B0=E8=87=82=E4=BD=BF=E8=83=BD=E6=9E=9A?= =?UTF-8?q?=E4=B8=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 FlowFailureCleanupService 用于清理设备端和后端媒体状态 - 添加 ARM_ENABLE 和 ARM_DISABLE 枚举值支持机械臂使能控制 - 在 TaskInstHolder 中集成失败清理逻辑 - 实现机器人录音、录像和系统状态的异常清理功能 - 添加 FlowFailureCleanupService 单元测试验证清理逻辑 - 注入 FlowFailureCleanupService 并在任务失败时执行清理操作 --- .../java/com/cmvr/test/enums/ActionEnum.java | 4 ++ .../test/flow/context/TaskInstHolder.java | 8 +++ .../control/FlowFailureCleanupService.java | 60 +++++++++++++++++++ .../FlowFailureCleanupServiceTest.java | 47 +++++++++++++++ 4 files changed, 119 insertions(+) create mode 100644 cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowFailureCleanupService.java create mode 100644 cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowFailureCleanupServiceTest.java diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java index 5b72d09..bdb72c3 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java @@ -65,6 +65,10 @@ public enum ActionEnum { // 末端运动 ARM_MOVE_TO_POINT("EDGE", "ARM_MOVE_TO_POINT", "运动到指定点"), 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", "开始说话"), diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java index 056a8ae..2522d41 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java @@ -4,12 +4,14 @@ import com.alibaba.fastjson2.JSONObject; import com.cmvr.test.enums.NodeTypeEnum; import com.cmvr.test.enums.TaskStatusEnum; 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.FlowExecutionEventPublisher; import com.cmvr.test.service.ITeNodeInstService; import com.cmvr.test.service.ITeTaskInstService; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import jakarta.annotation.Resource; import org.springframework.stereotype.Component; import java.util.List; @@ -22,6 +24,9 @@ import java.util.List; @RequiredArgsConstructor public class TaskInstHolder { + @Resource + private FlowFailureCleanupService flowFailureCleanupService; + private final ITeNodeInstService nodeInstService; private final TaskContextManager taskContextManager; private final ITeTaskInstService taskInstService; @@ -123,6 +128,9 @@ public class TaskInstHolder { // 终止和并发失败都可能先注销上下文,失败状态仍需落库。 TaskContext ctx = taskContextManager.get(instId); String moduleCode = ctx != null ? ctx.getModuleCode() : null; + if (ctx != null && flowFailureCleanupService != null) { + flowFailureCleanupService.cleanup(ctx); + } syncStatus(instId, TaskStatusEnum.FAILED); // 注销任务上下文 diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowFailureCleanupService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowFailureCleanupService.java new file mode 100644 index 0000000..423c699 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowFailureCleanupService.java @@ -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 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()); + } + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowFailureCleanupServiceTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowFailureCleanupServiceTest.java new file mode 100644 index 0000000..f396af2 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowFailureCleanupServiceTest.java @@ -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 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 proxy(Class type, String name, List 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; + }); + } +}