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; + }); + } +}