From 937426fcfc0e2ed842eecff849fe3bdb2b9f47e1 Mon Sep 17 00:00:00 2001 From: lixiaolong <702156524@qq.com> Date: Tue, 18 Aug 2026 14:43:43 +0800 Subject: [PATCH] =?UTF-8?q?feat(edge):=20=E6=96=B0=E5=A2=9EAGV=E5=B9=B3?= =?UTF-8?q?=E7=A7=BB=E5=92=8C=E6=9C=BA=E6=A2=B0=E8=87=82=E4=BD=BF=E8=83=BD?= =?UTF-8?q?=E6=8E=A7=E5=88=B6=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增AGV固定距离平移功能,支持指定速度和距离移动 - 新增机械臂使能开启/关闭控制,独立于移动操作 - 新增机械臂厂商扩展JSON指令执行接口 - 新增边缘系统安全状态查询和恢复功能 - 优化任务执行监听器,支持任务终态消息推送和错误处理 - 重构音频事件分类逻辑,支持预期标签匹配验证 - 新增任务日志状态管理和终端消息推送机制 --- .../web/controller/api/EdgeAgvController.java | 20 ++++++ .../web/controller/api/EdgeArmController.java | 9 +++ .../controller/api/EdgeSystemController.java | 15 ++++ .../listener/AimaFlowExecutionListener.java | 70 +++++++++++++++---- .../service/impl/AimaTestLogServiceImpl.java | 7 +- .../aima/support/AimaAnalysisVerdict.java | 36 ++++++++++ .../aima/support/AimaLogExecutionContext.java | 53 ++++++++++++++ .../AimaFlowExecutionListenerTest.java | 54 ++++++++++++++ .../aima/support/AimaAnalysisVerdictTest.java | 26 +++++++ .../support/AimaLogExecutionContextTest.java | 24 +++++++ .../cmvr-iot-grpc-client/pom.xml | 7 ++ .../client/model/agv/EdgeAgvTranslateVO.java | 27 +++++++ .../model/system/EdgeSafetyRecoveryVO.java | 25 +++++++ .../edge/client/service/EdgeAgvService.java | 5 ++ .../edge/client/service/EdgeArmService.java | 6 ++ .../client/service/EdgeSystemService.java | 6 ++ .../service/impl/EdgeAgvServiceImpl.java | 14 ++++ .../service/impl/EdgeArmServiceImpl.java | 12 ++++ .../service/impl/EdgeSystemServiceImpl.java | 70 +++++++++++++++++++ .../impl/EdgeSystemServiceImplTest.java | 41 +++++++++++ .../InspectionFlowExecutionListener.java | 32 ++++++--- .../java/com/cmvr/test/enums/ActionEnum.java | 2 +- .../interceptor/FlowBeforeInterceptor.java | 1 + .../operator/edge/EdgeArmOperateService.java | 16 +++-- .../llm/LLMMediaAnalysisOperateService.java | 36 +++++++++- .../service/FlowActionExecutorService.java | 10 ++- .../edge/EdgeArmPowerControlTest.java | 54 ++++++++++++++ .../LLMMediaAnalysisOperateServiceTest.java | 30 ++++++++ 28 files changed, 674 insertions(+), 34 deletions(-) create mode 100644 cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaAnalysisVerdict.java create mode 100644 cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java create mode 100644 cmvr-iot-aima/src/test/java/com/cmvr/aima/listener/AimaFlowExecutionListenerTest.java create mode 100644 cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaAnalysisVerdictTest.java create mode 100644 cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java create mode 100644 cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/agv/EdgeAgvTranslateVO.java create mode 100644 cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/system/EdgeSafetyRecoveryVO.java create mode 100644 cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/test/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImplTest.java create mode 100644 cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmPowerControlTest.java diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeAgvController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeAgvController.java index e8f20c0..e75a450 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeAgvController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeAgvController.java @@ -9,6 +9,7 @@ import com.cmvr.edge.client.model.agv.EdgeAgvNavigateToStationVO; import com.cmvr.edge.client.model.agv.EdgeAgvStartMappingVO; import com.cmvr.edge.client.model.agv.EdgeAgvUploadMapVO; import com.cmvr.edge.client.model.agv.EdgeAgvVelocityVO; +import com.cmvr.edge.client.model.agv.EdgeAgvTranslateVO; import com.cmvr.edge.client.service.EdgeAgvService; import io.swagger.annotations.Api; import io.swagger.annotations.ApiOperation; @@ -100,6 +101,25 @@ public class EdgeAgvController { return AjaxResult.ok(); } + @ApiOperation("固定距离平移") + @PostMapping("/translate") + public AjaxResult translate(@RequestBody @Validated EdgeAgvTranslateVO vo) { + AgvUtils.AgvTranslationMode mode = vo.getMode() == null + ? AgvUtils.AgvTranslationMode.AGV_TRANSLATION_MODE_ODOMETRY + : AgvUtils.AgvTranslationMode.forNumber(vo.getMode()); + if (mode == null) { + return AjaxResult.error("距离参考模式不合法,允许值:0=里程,1=定位"); + } + AgvUtils.AgvTranslation translation = AgvUtils.AgvTranslation.newBuilder() + .setDistance(vo.getDistance()) + .setVx(vo.getVx()) + .setVy(vo.getVy()) + .setMode(mode) + .build(); + edgeAgvService.translate(vo, translation); + return AjaxResult.ok(); + } + /** * 导航到指定坐标。 */ diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeArmController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeArmController.java index 7a5100a..b5695b6 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeArmController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeArmController.java @@ -4,11 +4,13 @@ import cmvr.api.ArmCommand; import com.cmvr.common.core.domain.AjaxResult; import com.cmvr.edge.client.model.EdgeCommonVO; import com.cmvr.edge.client.model.arm.*; +import com.cmvr.edge.client.model.system.EdgeSystemJsonCommandVO; import com.cmvr.edge.client.service.EdgeArmService; import io.swagger.annotations.Api; import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; import lombok.RequiredArgsConstructor; +import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; @@ -63,6 +65,13 @@ public class EdgeArmController { return AjaxResult.ok(); } + @ApiOperation("执行机械臂厂商扩展JSON指令") + @PostMapping("/executeJsonCommand") + public AjaxResult executeJsonCommand(@RequestBody @Validated EdgeSystemJsonCommandVO request) { + return AjaxResult.ok(com.alibaba.fastjson2.JSON.toJSONString( + edgeArmService.executeJsonCommand(request))); + } + @ApiOperation("关节空间运动") @PostMapping("/moveJ") public AjaxResult moveJ(@RequestBody EdgeArmMoveJVO moveJVO) { diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeSystemController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeSystemController.java index fdb528f..62482aa 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeSystemController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/api/EdgeSystemController.java @@ -6,6 +6,7 @@ import com.cmvr.common.core.domain.AjaxResult; import com.cmvr.edge.client.model.system.EdgeSystemJsonCommandVO; import com.cmvr.edge.client.model.system.EdgeSystemUpdateParamsVO; import com.cmvr.edge.client.model.system.EdgeActionQueueVO; +import com.cmvr.edge.client.model.system.EdgeSafetyRecoveryVO; import com.cmvr.edge.client.service.EdgeSystemService; import io.swagger.annotations.Api; import io.swagger.annotations.ApiOperation; @@ -16,6 +17,7 @@ import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; +import org.springframework.validation.annotation.Validated; @Api(tags = "边缘--系统功能") @RestController @@ -55,6 +57,19 @@ public class EdgeSystemController extends BaseController { return success(JSON.toJSONString(edgeSystemService.executeActionQueue(request))); } + @ApiOperation("Get edge safety state") + @GetMapping("/safetyState") + public AjaxResult safetyState(@RequestParam("robotId") String robotId) { + return success(JSON.toJSONString(edgeSystemService.getSafetyState(robotId))); + } + + @ApiOperation("Clear recoverable edge safety errors") + @PostMapping("/recoverSafetyState") + public AjaxResult recoverSafetyState(@RequestBody @Validated EdgeSafetyRecoveryVO request) { + return success(JSON.toJSONString(edgeSystemService.recoverSafetyState( + request.getRobotId(), request.getReason(), request.getTimeoutMs()))); + } + @ApiOperation("停止所有设备") @GetMapping("/stopAll") public AjaxResult stopAll(@RequestParam("robotId") String robotId) { diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java index c0a63f7..5d42466 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java @@ -10,6 +10,8 @@ import com.cmvr.aima.enums.AimaTaskStatusEnum; import com.cmvr.aima.service.IAimaTaskInstanceService; import com.cmvr.aima.service.IAimaTaskService; import com.cmvr.aima.service.IAimaTestLogService; +import com.cmvr.aima.support.AimaLogExecutionContext; +import com.cmvr.aima.support.AimaAnalysisVerdict; import com.cmvr.test.flow.runtime.event.FlowExecutionEvent; import com.cmvr.test.flow.runtime.event.FlowExecutionListener; import com.cmvr.test.flow.runtime.event.FlowExecutionModuleCodes; @@ -78,9 +80,7 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { } // 查询任务实例 - AimaTaskInstance taskInstance = aimaTaskInstanceService.lambdaQuery() - .eq(AimaTaskInstance::getTaskInsId, instId) - .one(); + AimaTaskInstance taskInstance = findTaskInstance(instId); if (taskInstance == null) { log.warn("未查询到爱玛测试实例,taskInsId:{}", instId); return; @@ -88,6 +88,7 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { int targetStatus = AimaTaskStatusEnum.RUNNING.getCode(); boolean needUpdateDb = false; + String terminalMessage = null; Date now = DateUtils.getNowDate(); switch (event.getEventType()) { @@ -95,12 +96,20 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { targetStatus = AimaTaskStatusEnum.SUCCESS.getCode(); needUpdateDb = true; clearCache(instId); - pushCompleteMsg(taskInstance.getId(), now, targetStatus); + terminalMessage = "任务执行完成"; break; case TASK_FAILED: targetStatus = AimaTaskStatusEnum.FAILED.getCode(); needUpdateDb = true; clearCache(instId); + terminalMessage = StrUtil.isBlank(event.getErrorMessage()) + ? "任务执行失败" : "任务执行失败:" + event.getErrorMessage(); + break; + case TASK_STOPPED: + targetStatus = AimaTaskStatusEnum.STOPPED.getCode(); + needUpdateDb = true; + clearCache(instId); + terminalMessage = "任务已终止"; break; case NODE_COMPLETED: default: @@ -112,6 +121,10 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { taskInstance.setEndTime(now); aimaTaskInstanceService.updateById(taskInstance); } + if (terminalMessage != null) { + pushTerminalMsg(taskInstance.getId(), now, targetStatus, terminalMessage, + event.getErrorMessage()); + } // 节点推送 - 记录日志 if (StrUtil.isNotBlank(nodeId)) { @@ -143,7 +156,12 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { progress = (itemTotalRatio + innerRatio) * 100; } - String info = event.getEventType() == FlowExecutionEvent.EventType.NODE_STARTED ? "开始执行" : "执行完成"; + String info = event.getEventType() == FlowExecutionEvent.EventType.NODE_STARTED + ? "开始执行" + : event.getEventType() == FlowExecutionEvent.EventType.NODE_FAILED + ? "执行失败" : "执行完成"; + int nodeLogStatus = event.getEventType() == FlowExecutionEvent.EventType.NODE_FAILED + ? AimaTaskStatusEnum.FAILED.getCode() : targetStatus; String logId = UUID.randomUUID().toString().replace("-", ""); // 四舍五入保留2位小数 progress = Math.round(progress * 100) / 100.0; @@ -166,14 +184,16 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { .logLevel("INFO") .logType("执行日志") .logContent("节点【" + event.getNodeName() + "】" + info) - .actualResult(outputParams == null || outputParams.isEmpty() - ? null : outputParams.toJSONString()) + .actualResult(AimaLogExecutionContext.attach(outputParams, itemId, + event.getItemOrder(), event.getItemOccurrence(), + event.getEventType().name(), event.getNodeType())) .screenshotUrl(imageUrl) .videoUrl(videoUrl) + .errorStack(event.getErrorMessage()) .logTime(now) - .status(targetStatus) + .status(nodeLogStatus) .executionStep(event.getNodeName()) - .testStatus("0") // 0未执行 + .testStatus(AimaAnalysisVerdict.resolve(event.getAction(), outputParams)) .durationMs(0L) .progress( progress) .build(); @@ -208,6 +228,24 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { return taskInsId + "|" + event.getItemId(); } + private AimaTaskInstance findTaskInstance(String taskInsId) { + for (int attempt = 0; attempt < 5; attempt++) { + AimaTaskInstance instance = aimaTaskInstanceService.lambdaQuery() + .eq(AimaTaskInstance::getTaskInsId, taskInsId) + .one(); + if (instance != null) { + return instance; + } + try { + Thread.sleep(50L); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + return null; + } + } + return null; + } + private String resolveTestCaseId(String taskInsId, String taskId, FlowExecutionEvent event) { if (event.getItemOrder() == null || event.getItemOrder() < 1) { return event.getItemId(); @@ -219,27 +257,29 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { } /** - * 推送完成任务信息 + * 保存并推送任务终态信息 */ - public void pushCompleteMsg(String dbInsId, Date now, int targetStatus) { + public void pushTerminalMsg(String dbInsId, Date now, int targetStatus, + String message, String errorMessage) { String logId = UUID.randomUUID().toString().replace("-", ""); - // 保存完成日志 + // 保存终态日志 AimaTestLog testLog = AimaTestLog.builder() .id(logId) .taskInstanceId(dbInsId) .logLevel("INFO") .logType("系统日志") - .logContent("任务执行完成") + .logContent(message) + .errorStack(errorMessage) .logTime(now) .status(targetStatus) .build(); try { aimaTestLogService.insertAimaTestLog(testLog); - log.info("任务完成日志保存成功,logId:{},taskInstanceId:{}", logId, dbInsId); + log.info("任务终态日志保存成功,logId:{},taskInstanceId:{},status:{}", logId, dbInsId, targetStatus); } catch (Exception e) { - log.error("任务完成日志保存失败,taskInstanceId:{}", dbInsId, e); + log.error("任务终态日志保存失败,taskInstanceId:{},status:{}", dbInsId, targetStatus, e); } try { messagePushService.pushToChannel("AimaTaskInstance", testLog); diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTestLogServiceImpl.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTestLogServiceImpl.java index 3107918..87887ad 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTestLogServiceImpl.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTestLogServiceImpl.java @@ -10,6 +10,7 @@ import com.cmvr.common.utils.SecurityUtils; import com.cmvr.aima.mapper.AimaTestLogMapper; import com.cmvr.aima.domain.AimaTestLog; import com.cmvr.aima.service.IAimaTestLogService; +import com.cmvr.aima.support.AimaLogExecutionContext; import org.springframework.transaction.annotation.Transactional; /** @@ -26,12 +27,14 @@ public class AimaTestLogServiceImpl extends ServiceImpl selectAimaTestLogList(AimaTestLog aimaTestLog) { - return aimaTestLogMapper.selectAimaTestLogList(aimaTestLog); + List logs = aimaTestLogMapper.selectAimaTestLogList(aimaTestLog); + logs.forEach(AimaLogExecutionContext::hydrate); + return logs; } @Override diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaAnalysisVerdict.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaAnalysisVerdict.java new file mode 100644 index 0000000..65d8c8c --- /dev/null +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaAnalysisVerdict.java @@ -0,0 +1,36 @@ +package com.cmvr.aima.support; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.test.enums.ActionEnum; + +/** + * Converts media-analysis output to the existing Aima test_status protocol. + */ +public final class AimaAnalysisVerdict { + + public static final String PASSED = "1"; + public static final String NOT_PASSED = "2"; + + private AimaAnalysisVerdict() { + } + + public static String resolve(ActionEnum action, JSONObject output) { + if (!isAnalysisAction(action) || output == null || output.isEmpty()) { + return null; + } + String analysisStatus = output.getString("analysisStatus"); + if ("FAILED".equalsIgnoreCase(analysisStatus) || "ERROR".equalsIgnoreCase(analysisStatus)) { + return NOT_PASSED; + } + Boolean passed = output.getBoolean("passed"); + if (passed != null) { + return passed ? PASSED : NOT_PASSED; + } + Boolean matched = output.getBoolean("matched"); + return Boolean.TRUE.equals(matched) ? PASSED : NOT_PASSED; + } + + private static boolean isAnalysisAction(ActionEnum action) { + return action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE; + } +} diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java new file mode 100644 index 0000000..64c0305 --- /dev/null +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java @@ -0,0 +1,53 @@ +package com.cmvr.aima.support; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.aima.domain.AimaTestLog; + +/** + * Persists flow item identity inside the existing actual_result column. + */ +public final class AimaLogExecutionContext { + + private static final String CONTEXT_KEY = "_executionContext"; + + private AimaLogExecutionContext() { + } + + public static String attach(JSONObject outputParams, String itemId, + Integer itemOrder, Integer itemOccurrence) { + return attach(outputParams, itemId, itemOrder, itemOccurrence, null, null); + } + + public static String attach(JSONObject outputParams, String itemId, + Integer itemOrder, Integer itemOccurrence, + String eventType, String nodeType) { + JSONObject result = outputParams == null + ? new JSONObject() + : JSONObject.parseObject(outputParams.toJSONString()); + JSONObject context = new JSONObject(); + context.put("itemId", itemId); + context.put("itemOrder", itemOrder); + context.put("itemOccurrence", itemOccurrence); + context.put("eventType", eventType); + context.put("nodeType", nodeType); + result.put(CONTEXT_KEY, context); + return result.toJSONString(); + } + + public static AimaTestLog hydrate(AimaTestLog log) { + if (log == null || log.getActualResult() == null || log.getActualResult().isBlank()) { + return log; + } + try { + JSONObject context = JSONObject.parseObject(log.getActualResult()).getJSONObject(CONTEXT_KEY); + if (context != null) { + log.setItemId(context.getString("itemId")); + log.setItemOrder(context.getInteger("itemOrder")); + log.setItemOccurrence(context.getInteger("itemOccurrence")); + } + } catch (RuntimeException ignored) { + // Historical actual_result values are not guaranteed to be JSON. + } + return log; + } +} diff --git a/cmvr-iot-aima/src/test/java/com/cmvr/aima/listener/AimaFlowExecutionListenerTest.java b/cmvr-iot-aima/src/test/java/com/cmvr/aima/listener/AimaFlowExecutionListenerTest.java new file mode 100644 index 0000000..807f76a --- /dev/null +++ b/cmvr-iot-aima/src/test/java/com/cmvr/aima/listener/AimaFlowExecutionListenerTest.java @@ -0,0 +1,54 @@ +package com.cmvr.aima.listener; + +import com.cmvr.aima.domain.AimaTestLog; +import com.cmvr.aima.service.IAimaTestLogService; +import com.cmvr.framework.websocket.service.MessagePushService; +import org.junit.Test; + +import java.lang.reflect.Proxy; +import java.util.Date; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertSame; + +public class AimaFlowExecutionListenerTest { + + @Test + public void failedTerminalLogIsPersistedAndPushed() { + AtomicReference inserted = new AtomicReference<>(); + IAimaTestLogService logService = (IAimaTestLogService) Proxy.newProxyInstance( + IAimaTestLogService.class.getClassLoader(), + new Class[]{IAimaTestLogService.class}, + (proxy, method, args) -> { + if ("insertAimaTestLog".equals(method.getName())) { + inserted.set((AimaTestLog) args[0]); + return 1; + } + return null; + }); + CapturingMessagePushService messagePushService = new CapturingMessagePushService(); + AimaFlowExecutionListener listener = new AimaFlowExecutionListener( + null, logService, null, messagePushService); + + listener.pushTerminalMsg("instance-1", new Date(), 2, + "任务执行失败:机械臂异常", "机械臂异常"); + + assertEquals(Integer.valueOf(2), inserted.get().getStatus()); + assertEquals("任务执行失败:机械臂异常", inserted.get().getLogContent()); + assertEquals("机械臂异常", inserted.get().getErrorStack()); + assertEquals("AimaTaskInstance", messagePushService.channel); + assertSame(inserted.get(), messagePushService.payload); + } + + private static class CapturingMessagePushService extends MessagePushService { + private String channel; + private Object payload; + + @Override + public void pushToChannel(String channel, Object payload) { + this.channel = channel; + this.payload = payload; + } + } +} diff --git a/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaAnalysisVerdictTest.java b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaAnalysisVerdictTest.java new file mode 100644 index 0000000..d0a085f --- /dev/null +++ b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaAnalysisVerdictTest.java @@ -0,0 +1,26 @@ +package com.cmvr.aima.support; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.test.enums.ActionEnum; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; + +public class AimaAnalysisVerdictTest { + + @Test + public void resolvesOnlyAnalysisNodeResults() { + assertEquals(AimaAnalysisVerdict.PASSED, + AimaAnalysisVerdict.resolve(ActionEnum.VIDEO_ANALYZE, + new JSONObject().fluentPut("passed", true))); + assertEquals(AimaAnalysisVerdict.NOT_PASSED, + AimaAnalysisVerdict.resolve(ActionEnum.AUDIO_EVENT_CLASSIFY, + new JSONObject().fluentPut("matched", false))); + assertEquals(AimaAnalysisVerdict.NOT_PASSED, + AimaAnalysisVerdict.resolve(ActionEnum.VIDEO_ANALYZE, + new JSONObject().fluentPut("analysisStatus", "FAILED"))); + assertNull(AimaAnalysisVerdict.resolve(ActionEnum.SLEEP, + new JSONObject().fluentPut("passed", true))); + } +} diff --git a/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java new file mode 100644 index 0000000..6e2b9e9 --- /dev/null +++ b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java @@ -0,0 +1,24 @@ +package com.cmvr.aima.support; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.aima.domain.AimaTestLog; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; + +public class AimaLogExecutionContextTest { + + @Test + public void preservesOutputAndHydratesRepeatedItemIdentity() { + JSONObject output = new JSONObject().fluentPut("videoUrl", "/media/result.mp4"); + String actualResult = AimaLogExecutionContext.attach(output, "case-1", 4, 2); + AimaTestLog log = AimaTestLog.builder().actualResult(actualResult).build(); + + AimaLogExecutionContext.hydrate(log); + + assertEquals("/media/result.mp4", JSONObject.parseObject(actualResult).getString("videoUrl")); + assertEquals("case-1", log.getItemId()); + assertEquals(Integer.valueOf(4), log.getItemOrder()); + assertEquals(Integer.valueOf(2), log.getItemOccurrence()); + } +} diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/pom.xml b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/pom.xml index 0d87ffb..8949143 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/pom.xml +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/pom.xml @@ -64,6 +64,13 @@ 0.14.0 + + junit + junit + 4.13.2 + test + + diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/agv/EdgeAgvTranslateVO.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/agv/EdgeAgvTranslateVO.java new file mode 100644 index 0000000..6a273d0 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/agv/EdgeAgvTranslateVO.java @@ -0,0 +1,27 @@ +package com.cmvr.edge.client.model.agv; + +import com.cmvr.edge.client.model.EdgeCommonVO; +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import jakarta.validation.constraints.DecimalMin; +import lombok.Data; +import lombok.EqualsAndHashCode; + +@Data +@EqualsAndHashCode(callSuper = true) +@ApiModel("AGV fixed-distance translation request") +public class EdgeAgvTranslateVO extends EdgeCommonVO { + + @DecimalMin(value = "0.0", inclusive = false, message = "平移距离必须大于0") + @ApiModelProperty(value = "平移距离绝对值,单位:米", required = true) + private double distance; + + @ApiModelProperty(value = "车体X方向速度,单位:米/秒", required = true) + private double vx; + + @ApiModelProperty(value = "车体Y方向速度,单位:米/秒", required = true) + private double vy; + + @ApiModelProperty("距离参考模式:0=里程,1=定位") + private Integer mode; +} diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/system/EdgeSafetyRecoveryVO.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/system/EdgeSafetyRecoveryVO.java new file mode 100644 index 0000000..2d83356 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/model/system/EdgeSafetyRecoveryVO.java @@ -0,0 +1,25 @@ +package com.cmvr.edge.client.model.system; + +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import jakarta.validation.constraints.Max; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotBlank; +import lombok.Data; + +@Data +@ApiModel("Edge safety recovery request") +public class EdgeSafetyRecoveryVO { + + @NotBlank(message = "机器人ID不能为空") + @ApiModelProperty(value = "机器人ID", required = true) + private String robotId; + + @ApiModelProperty("恢复原因") + private String reason; + + @Min(value = 1000, message = "超时时间不能小于1000毫秒") + @Max(value = 120000, message = "超时时间不能大于120000毫秒") + @ApiModelProperty("超时时间,单位毫秒") + private Integer timeoutMs; +} diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeAgvService.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeAgvService.java index aff3095..f027152 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeAgvService.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeAgvService.java @@ -45,6 +45,11 @@ public interface EdgeAgvService { */ void clearFault(EdgeCommonVO edgeCommonVO); + /** + * 按指定速度平移固定距离。 + */ + void translate(EdgeCommonVO edgeCommonVO, AgvUtils.AgvTranslation translation); + /** * 导航到指定位置 * diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeArmService.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeArmService.java index bba3760..a80378d 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeArmService.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeArmService.java @@ -3,6 +3,7 @@ package com.cmvr.edge.client.service; import cmvr.api.ArmCommand; import cmvr.api.Common; import com.cmvr.edge.client.model.EdgeCommonVO; +import com.cmvr.edge.client.model.system.EdgeSystemJsonCommandVO; /** * 边缘系统机械臂服务 @@ -38,6 +39,11 @@ public interface EdgeArmService { */ Common.CommandHeader.Feedback clearFault(EdgeCommonVO edgeCommonVO); + /** + * 执行机械臂厂商扩展 JSON 指令。 + */ + Common.JsonDeviceCommand.Feedback executeJsonCommand(EdgeSystemJsonCommandVO request); + /** * 关节空间运动(命令码:moveJ) * 控制机械臂通过关节空间插值运动到目标关节位置 diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeSystemService.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeSystemService.java index 121f7d1..7dc4f5e 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeSystemService.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeSystemService.java @@ -1,6 +1,7 @@ package com.cmvr.edge.client.service; import cmvr.api.Common; +import cmvr.api.SafetyCommand; import cmvr.api.SystemCommand; import com.cmvr.device.domain.DeDeviceRegistration; import com.cmvr.edge.client.model.system.EdgeSystemJsonCommandVO; @@ -27,6 +28,11 @@ public interface EdgeSystemService { SystemCommand.ActionQueueCommand.Feedback executeActionQueue(EdgeActionQueueVO request); + SafetyCommand.GetSafetyStateCommand.Feedback getSafetyState(String robotId); + + SafetyCommand.RecoverSafetyStateCommand.Feedback recoverSafetyState( + String robotId, String reason, Integer timeoutMs); + public List deviceList(String robotId); diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeAgvServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeAgvServiceImpl.java index 4541263..b87640e 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeAgvServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeAgvServiceImpl.java @@ -65,6 +65,20 @@ public class EdgeAgvServiceImpl implements EdgeAgvService { executeGrpcCall(() -> stub.clearFault(request)); } + @Override + public void translate(EdgeCommonVO edgeCommonVO, AgvUtils.AgvTranslation translation) { + if (translation == null || translation.getDistance() <= 0) { + throw new GlobalException("平移距离必须大于0"); + } + AgvServiceGrpc.AgvServiceBlockingStub stub = grpcServiceManager.getGrpcClient( + edgeCommonVO.getRobotId(), AgvServiceGrpc.AgvServiceBlockingStub.class); + AgvCommand.AgvTranslateCommand.Request request = AgvCommand.AgvTranslateCommand.Request.newBuilder() + .setHeader(EdgeCommonUtil.buildRequest(edgeCommonVO.getDeviceId())) + .setTranslation(translation) + .build(); + executeGrpcCall(() -> stub.translate(request)); + } + @Override public void navigateToPose(EdgeCommonVO edgeCommonVO, AgvUtils.AgvPose2d pose) { if (pose == null) { diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeArmServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeArmServiceImpl.java index 8f2b75d..63a403d 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeArmServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeArmServiceImpl.java @@ -7,6 +7,7 @@ import cn.hutool.core.util.StrUtil; import com.cmvr.common.exception.GlobalException; import com.cmvr.edge.client.manage.GrpcServiceManager; import com.cmvr.edge.client.model.EdgeCommonVO; +import com.cmvr.edge.client.model.system.EdgeSystemJsonCommandVO; import com.cmvr.edge.client.service.EdgeArmService; import com.cmvr.edge.client.utils.EdgeCommonUtil; import lombok.RequiredArgsConstructor; @@ -52,6 +53,17 @@ public class EdgeArmServiceImpl implements EdgeArmService { return executeGrpcCall(() -> stub.clearFault(request)); } + @Override + public Common.JsonDeviceCommand.Feedback executeJsonCommand(EdgeSystemJsonCommandVO request) { + ArmServiceGrpc.ArmServiceBlockingStub stub = grpcServiceManager.getGrpcClient( + request.getRobotId(), ArmServiceGrpc.ArmServiceBlockingStub.class); + Common.JsonDeviceCommand.Request grpcRequest = Common.JsonDeviceCommand.Request.newBuilder() + .setHeader(EdgeCommonUtil.buildRequest(request.getDeviceId())) + .setRequestJson(request.getRequestJson() == null ? "" : request.getRequestJson()) + .build(); + return executeGrpcCall(() -> stub.executeJsonCommand(grpcRequest)); + } + @Override public ArmCommand.MoveJ.Response moveJ(EdgeCommonVO edgeCommonVO, double[] target, Double velocity, Double acceleration, Double blendRadius, 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 39bf616..a512f35 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 @@ -7,6 +7,7 @@ import cmvr.api.AgvCommand; import cmvr.msgs.AgvUtils; import cmvr.api.SystemCommand; import cmvr.api.SystemServiceGrpc; +import cmvr.api.SafetyCommand; import com.alibaba.fastjson2.JSON; import com.cmvr.device.domain.DeDeviceRegistration; import com.cmvr.edge.client.manage.GrpcServiceManager; @@ -116,6 +117,75 @@ public class EdgeSystemServiceImpl implements EdgeSystemService { return stub.withDeadlineAfter(deadlineMs, TimeUnit.MILLISECONDS).executeActionQueue(builder.build()); } + @Override + public SafetyCommand.GetSafetyStateCommand.Feedback getSafetyState(String robotId) { + if (robotId == null || robotId.trim().isEmpty()) { + throw new GlobalException("机器人ID不能为空"); + } + SystemServiceGrpc.SystemServiceBlockingStub stub = grpcServiceManager.getGrpcClient( + robotId.trim(), SystemServiceGrpc.SystemServiceBlockingStub.class); + SafetyCommand.GetSafetyStateCommand.Request request = + SafetyCommand.GetSafetyStateCommand.Request.newBuilder() + .setScope(allDevicesScope()) + .build(); + SafetyCommand.GetSafetyStateCommand.Feedback feedback = stub.getSafetyState(request); + requireSuccessfulHeader(feedback == null || !feedback.hasHeader() ? null : feedback.getHeader(), + "获取机器人安全状态"); + return feedback; + } + + @Override + public SafetyCommand.RecoverSafetyStateCommand.Feedback recoverSafetyState( + String robotId, String reason, Integer timeoutMs) { + SafetyCommand.GetSafetyStateCommand.Feedback current = getSafetyState(robotId); + SafetyCommand.RecoverSafetyStateCommand.Request request = buildSafetyRecoveryRequest( + current.getSafetyEpoch(), reason, timeoutMs); + SystemServiceGrpc.SystemServiceBlockingStub stub = grpcServiceManager.getGrpcClient( + robotId.trim(), SystemServiceGrpc.SystemServiceBlockingStub.class); + SafetyCommand.RecoverSafetyStateCommand.Feedback feedback = stub + .withDeadlineAfter(request.getTimeoutMs() + 5000L, TimeUnit.MILLISECONDS) + .recoverSafetyState(request); + requireSuccessfulHeader(feedback == null || !feedback.hasHeader() ? null : feedback.getHeader(), + "恢复机器人安全状态"); + if (!isSuccessfulRecoveryResult(feedback.getResult())) { + throw new GlobalException("恢复机器人安全状态未完成: " + feedback.getResult().name()); + } + return feedback; + } + + SafetyCommand.RecoverSafetyStateCommand.Request buildSafetyRecoveryRequest( + long expectedSafetyEpoch, String reason, Integer timeoutMs) { + int effectiveTimeoutMs = timeoutMs == null ? 15000 : Math.max(1000, Math.min(timeoutMs, 120000)); + return SafetyCommand.RecoverSafetyStateCommand.Request.newBuilder() + .setRecoveryId(UUID.randomUUID().toString()) + .setScope(allDevicesScope()) + .setExpectedSafetyEpoch(expectedSafetyEpoch) + .setMode(SafetyCommand.RecoverSafetyStateCommand.Mode.CLEAR_SOFTWARE_LATCH) + .setReason(reason == null || reason.trim().isEmpty() + ? "Platform manual safety recovery" : reason.trim()) + .setTimeoutMs(effectiveTimeoutMs) + .build(); + } + + private boolean isSuccessfulRecoveryResult(SafetyCommand.SafetyOperationResult result) { + return result == SafetyCommand.SafetyOperationResult.SAFETY_OPERATION_RESULT_SUCCEEDED + || result == SafetyCommand.SafetyOperationResult.SAFETY_OPERATION_RESULT_RECOVERED + || result == SafetyCommand.SafetyOperationResult.SAFETY_OPERATION_RESULT_NOTHING_TO_RECOVER; + } + + private SafetyCommand.SafetyScope allDevicesScope() { + return SafetyCommand.SafetyScope.newBuilder().setAllDevices(true).build(); + } + + private void requireSuccessfulHeader(Common.CommandHeader.Feedback header, String operation) { + if (header == null) { + throw new GlobalException(operation + "失败: 响应为空"); + } + if (!header.getSuccess()) { + throw new GlobalException(operation + "失败: " + header.getErrorMessage()); + } + } + private SystemCommand.ActionStep buildActionStep(EdgeActionQueueVO.Step step, String stepId) { if (step == null || step.getType() == null) { throw new GlobalException("Action queue step type cannot be empty"); diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/test/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImplTest.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/test/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImplTest.java new file mode 100644 index 0000000..e8b7314 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/test/java/com/cmvr/edge/client/service/impl/EdgeSystemServiceImplTest.java @@ -0,0 +1,41 @@ +package com.cmvr.edge.client.service.impl; + +import cmvr.api.SafetyCommand; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class EdgeSystemServiceImplTest { + + private final EdgeSystemServiceImpl service = new EdgeSystemServiceImpl(null); + + @Test + public void buildsAllDeviceSoftwareLatchRecoveryRequest() { + SafetyCommand.RecoverSafetyStateCommand.Request request = + service.buildSafetyRecoveryRequest(27L, " manual recovery ", 3000); + + assertFalse(request.getRecoveryId().isBlank()); + assertTrue(request.getScope().getAllDevices()); + assertEquals(27L, request.getExpectedSafetyEpoch()); + assertEquals(SafetyCommand.RecoverSafetyStateCommand.Mode.CLEAR_SOFTWARE_LATCH, request.getMode()); + assertEquals("manual recovery", request.getReason()); + assertEquals(3000, request.getTimeoutMs()); + } + + @Test + public void appliesRecoveryDefaultsAndTimeoutBounds() { + SafetyCommand.RecoverSafetyStateCommand.Request defaults = + service.buildSafetyRecoveryRequest(0L, " ", null); + SafetyCommand.RecoverSafetyStateCommand.Request minimum = + service.buildSafetyRecoveryRequest(0L, "test", 100); + SafetyCommand.RecoverSafetyStateCommand.Request maximum = + service.buildSafetyRecoveryRequest(0L, "test", 200000); + + assertEquals("Platform manual safety recovery", defaults.getReason()); + assertEquals(15000, defaults.getTimeoutMs()); + assertEquals(1000, minimum.getTimeoutMs()); + assertEquals(120000, maximum.getTimeoutMs()); + } +} diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java index f2d08b0..f837714 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java @@ -91,18 +91,24 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { int targetStatus = TaskStatusEnum.RUNNING.getCode(); boolean needUpdateDb = false; + String terminalMessage = null; Date now = DateUtils.getNowDate(); switch (event.getEventType()) { case TASK_COMPLETED: targetStatus = TaskStatusEnum.SUCCESS.getCode(); needUpdateDb = true; - pushCompleteMsg(taskInstance.getId(), now, targetStatus); + terminalMessage = "任务执行完成"; break; case TASK_FAILED: targetStatus = TaskStatusEnum.FAILED.getCode(); needUpdateDb = true; + terminalMessage = StrUtil.isBlank(event.getErrorMessage()) + ? "任务执行失败" : "任务执行失败:" + event.getErrorMessage(); break; case TASK_STOPPED: + targetStatus = TaskStatusEnum.STOPPED.getCode(); + needUpdateDb = true; + terminalMessage = "任务已终止"; break; case NODE_COMPLETED: default: @@ -114,6 +120,9 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { taskInstance.setEndTime(now); inspectionTaskInstanceService.updateById(taskInstance); } + if (terminalMessage != null) { + pushTerminalMsg(taskInstance.getId(), now, targetStatus, terminalMessage); + } if (FlowExecutionEvent.EventType.NODE_COMPLETED == event.getEventType()) { try { @@ -155,7 +164,12 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { progress = (itemTotalRatio + innerRatio) * 100; } - String info = event.getEventType() == FlowExecutionEvent.EventType.NODE_STARTED ? "开始执行" : "执行完成"; + String info = event.getEventType() == FlowExecutionEvent.EventType.NODE_STARTED + ? "开始执行" + : event.getEventType() == FlowExecutionEvent.EventType.NODE_FAILED + ? "执行失败" : "执行完成"; + int nodeLogStatus = event.getEventType() == FlowExecutionEvent.EventType.NODE_FAILED + ? TaskStatusEnum.FAILED.getCode() : targetStatus; String logId = UUID.randomUUID().toString().replace("-", ""); // 四舍五入保留2位小数 progress = Math.round(progress * 100) / 100.0; @@ -190,7 +204,7 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { .mediaUrl(mediaUrl) .extraInfo(outputParams == null || outputParams.isEmpty() ? null : outputParams.toJSONString()) - .status(targetStatus) + .status(nodeLogStatus) .progress(progress) .logTime(now) .build(); @@ -247,26 +261,26 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { } /** - * 推送完成任务信息(参数dbInsId为数据库主键ID) + * 保存并推送任务终态信息(参数dbInsId为数据库主键ID) */ - public void pushCompleteMsg(String dbInsId, Date now, int targetStatus) { + public void pushTerminalMsg(String dbInsId, Date now, int targetStatus, String message) { String logId = UUID.randomUUID().toString().replace("-", ""); - // 保存完成日志 + // 保存终态日志 InspectionTaskLog taskLog = InspectionTaskLog.builder() .id(logId) .taskInstanceId(dbInsId) .logType(InspectionLogTypeEnum.TEXT.getCode()) - .logContent("任务执行完成") + .logContent(message) .status(targetStatus) .logTime(now) .build(); try { inspectionTaskLogService.insertInspectionTaskLog(taskLog); - log.info("任务完成日志保存成功,logId:{},taskInstanceId:{}", logId, dbInsId); + log.info("任务终态日志保存成功,logId:{},taskInstanceId:{},status:{}", logId, dbInsId, targetStatus); } catch (Exception e) { - log.error("任务完成日志保存失败,taskInstanceId:{}", dbInsId, e); + log.error("任务终态日志保存失败,taskInstanceId:{},status:{}", dbInsId, targetStatus, e); } try { messagePushService.pushToChannel("InspectionTaskInstance", taskLog); 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 bdb72c3..83fde61 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 @@ -88,7 +88,7 @@ public enum ActionEnum { INSPECTION_METER_RECOGNIZE("LLM", "INSPECTION_METER_RECOGNIZE", "巡检仪表读数识别"), AI_TTS("LLM", "AI_TTS", "tts语音播放"), GET_CURRENT_PAGE("LLM", "GET_CURRENT_PAGE", "获取当前页面名称"), - AUDIO_EVENT_CLASSIFY("LLM", "AUDIO_EVENT_CLASSIFY", "声音事件识别"), + AUDIO_EVENT_CLASSIFY("LLM", "AUDIO_EVENT_CLASSIFY", "声音类型检测"), VIDEO_ANALYZE("LLM", "VIDEO_ANALYZE", "视频智能分析"), // 触控交互 diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java index bbea189..cdb5d88 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java @@ -46,6 +46,7 @@ public class FlowBeforeInterceptor extends AbstractFlowMsgPreInterceptor{ .nodeCount(message.getGraph().allNodeIds().size()) .nodeName(message.getNodeName()) .nodeType(message.getNodeType()) + .action(message.getAction()) .build(); eventPublisher.publishEvent(event); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmOperateService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmOperateService.java index ab7159b..9a65fd4 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmOperateService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmOperateService.java @@ -13,8 +13,6 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; -import java.util.Objects; - /** * 机械臂边缘操作服务 * 处理机械臂相关的设备行为,如末端运动到指定位置 @@ -51,10 +49,18 @@ public class EdgeArmOperateService implements EdgeOperateService { edgeCommonVO.setRobotId(robotId); edgeCommonVO.setDeviceId(deviceId); - // 先调用开启力矩 - edgeArmService.torqueOn(edgeCommonVO); + // Power state changes are explicit workflow actions; move actions never enable automatically. + if (action == null) { + throw new GlobalException("机械臂动作不能为空"); + } - if (Objects.requireNonNull(action) == ActionEnum.ARM_MOVE_TO_POINT) { + if (action == ActionEnum.ARM_ENABLE) { + log.info("开启机械臂使能,设备ID: {}", deviceId); + edgeArmService.torqueOn(edgeCommonVO); + } else if (action == ActionEnum.ARM_DISABLE) { + log.info("关闭机械臂使能,设备ID: {}", deviceId); + edgeArmService.torqueOff(edgeCommonVO); + } else if (action == ActionEnum.ARM_MOVE_TO_POINT) { log.info("执行机械臂末端运动到指定点操作,设备ID: {}", deviceId); // 从输入参数中获取坐标信息 diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java index ecedb39..7bf2587 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java @@ -49,7 +49,16 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService { context.put("trial", message.isTrial()); JSONObject options = new JSONObject(); - if (!audio) { + String expectedLabel = null; + String expectedLabelName = null; + if (audio) { + expectedLabel = StringUtils.trimToNull(input.getString("expectedLabel")); + expectedLabelName = StringUtils.trimToNull(input.getString("expectedLabelName")); + if (expectedLabel != null) { + options.put("expectedLabel", expectedLabel); + options.put("expectedLabelName", StringUtils.defaultIfBlank(expectedLabelName, expectedLabel)); + } + } else { options.put("instruction", StringUtils.defaultString(input.getString("instruction"))); String analysisMode = StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO") .trim().toUpperCase(); @@ -85,9 +94,34 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService { output.put("analysisModel", response.getJSONObject("model")); output.put("analysisTimingMs", response.getLong("timingMs")); output.put("mediaUrl", mediaUrl); + if (audio) { + applyAudioVerdict(output, expectedLabel, expectedLabelName); + } return TaskNodeExecuteResult.success(output); } + private void applyAudioVerdict(JSONObject output, String expectedLabel, String expectedLabelName) { + boolean classified = Boolean.TRUE.equals(output.getBoolean("matched")); + if (expectedLabel == null) { + // Existing workflows did not define a target type. Preserve their old + // "recognized any known sound" semantics until they are edited and saved. + output.put("passed", classified); + output.put("decisionReason", classified ? "已识别到已知声音类型" : "未识别到已知声音类型"); + return; + } + String actualLabel = StringUtils.trimToEmpty(output.getString("label")); + String actualLabelName = StringUtils.trimToEmpty(output.getString("labelName")); + boolean passed = classified && (expectedLabel.equalsIgnoreCase(actualLabel) + || expectedLabel.equalsIgnoreCase(actualLabelName)); + String displayName = StringUtils.defaultIfBlank(expectedLabelName, expectedLabel); + output.put("expectedLabel", expectedLabel); + output.put("expectedLabelName", displayName); + output.put("passed", passed); + output.put("decisionReason", passed + ? "检测到目标声音:" + displayName + : "未检测到目标声音:" + displayName); + } + private JSONObject buildVideoTuning(JSONObject source) { JSONObject tuning = new JSONObject(); if (source == null || source.isEmpty()) { diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java index c8da129..1315339 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java @@ -198,10 +198,18 @@ public class FlowActionExecutorService { // ==== 机械臂 ==== case ARM_MOVE_TO_POINT: - case ARM_MOVE_TO_J: { + case ARM_MOVE_TO_J: + case ARM_ENABLE: + case ARM_DISABLE: { TaskNodeExecuteMessage message = buildSingleNodeMessage(action, payload); message.setRobotId(robotId); edgeArmOperateService.execute(message); + if (action == ActionEnum.ARM_ENABLE) { + return "机械臂使能开启成功"; + } + if (action == ActionEnum.ARM_DISABLE) { + return "机械臂使能关闭成功"; + } return action == ActionEnum.ARM_MOVE_TO_J ? "机械臂关节空间运动任务下发成功" : "机械臂末端运动任务下发成功"; diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmPowerControlTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmPowerControlTest.java new file mode 100644 index 0000000..4b150de --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/EdgeArmPowerControlTest.java @@ -0,0 +1,54 @@ +package com.cmvr.test.flow.runtime.operator.edge; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.edge.client.service.EdgeArmService; +import com.cmvr.test.enums.ActionEnum; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +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 EdgeArmPowerControlTest { + + @Test + public void powerIsExplicitAndMoveDoesNotEnableAutomatically() { + List methods = new ArrayList<>(); + EdgeArmService edgeArmService = (EdgeArmService) Proxy.newProxyInstance( + EdgeArmService.class.getClassLoader(), + new Class[]{EdgeArmService.class}, + (proxy, method, args) -> { + methods.add(method.getName()); + return null; + }); + EdgeArmOperateService service = new EdgeArmOperateService(edgeArmService); + + assertTrue(service.execute(message(ActionEnum.ARM_ENABLE, baseParams())).isSuccess()); + assertTrue(service.execute(message(ActionEnum.ARM_MOVE_TO_POINT, moveParams())).isSuccess()); + assertTrue(service.execute(message(ActionEnum.ARM_DISABLE, baseParams())).isSuccess()); + + assertEquals(List.of("torqueOn", "moveL", "torqueOff"), methods); + } + + private TaskNodeExecuteMessage message(ActionEnum action, JSONObject input) { + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setAction(action); + message.setRobotId("robot-1"); + message.setInputParams(input); + return message; + } + + private JSONObject baseParams() { + return new JSONObject().fluentPut("deviceId", "arm-1"); + } + + private JSONObject moveParams() { + return baseParams() + .fluentPut("x", 0.1).fluentPut("y", 0.2).fluentPut("z", 0.3) + .fluentPut("rx", 0.4).fluentPut("ry", 0.5).fluentPut("rz", 0.6); + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java index 455770a..9a1f8b6 100644 --- a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java @@ -45,6 +45,36 @@ public class LLMMediaAnalysisOperateServiceTest { assertEquals("https://files.example/latest.wav", captured.get().getString("mediaUrl")); assertEquals("AUDIO_CLASSIFICATION", captured.get().getString("analysisType")); assertEquals("SUCCEEDED", result.getOutputParams().getString("analysisStatus")); + assertTrue(result.getOutputParams().getBooleanValue("passed")); + } + + @Test + public void comparesRecognizedAudioWithExpectedType() { + AtomicReference captured = new AtomicReference<>(); + MediaAnalysisClient client = request -> { + captured.set(request); + return new JSONObject() + .fluentPut("status", "SUCCEEDED") + .fluentPut("analysisType", "AUDIO_CLASSIFICATION") + .fluentPut("profileCode", "aima.power_state.v1") + .fluentPut("result", new JSONObject() + .fluentPut("label", "POWER_OFF") + .fluentPut("labelName", "关机") + .fluentPut("matched", true)); + }; + LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(client); + JSONObject input = new JSONObject() + .fluentPut("profileCode", "aima.power_state.v1") + .fluentPut("expectedLabel", "POWER_ON") + .fluentPut("expectedLabelName", "开机") + .fluentPut("audioUrl", "https://files.example/audio.wav"); + + TaskNodeExecuteResult result = service.execute(message(ActionEnum.AUDIO_EVENT_CLASSIFY, input)); + + assertEquals("POWER_ON", captured.get().getJSONObject("options").getString("expectedLabel")); + assertEquals("POWER_ON", result.getOutputParams().getString("expectedLabel")); + assertEquals("开机", result.getOutputParams().getString("expectedLabelName")); + assertEquals(false, result.getOutputParams().getBooleanValue("passed")); } @Test