feat(edge): 新增AGV平移和机械臂使能控制功能

- 新增AGV固定距离平移功能,支持指定速度和距离移动
- 新增机械臂使能开启/关闭控制,独立于移动操作
- 新增机械臂厂商扩展JSON指令执行接口
- 新增边缘系统安全状态查询和恢复功能
- 优化任务执行监听器,支持任务终态消息推送和错误处理
- 重构音频事件分类逻辑,支持预期标签匹配验证
- 新增任务日志状态管理和终端消息推送机制
This commit is contained in:
lixiaolong 2026-08-18 14:43:43 +08:00
parent 0b05438712
commit 937426fcfc
28 changed files with 674 additions and 34 deletions

View File

@ -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();
}
/**
* 导航到指定坐标。
*/

View File

@ -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) {

View File

@ -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) {

View File

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

View File

@ -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<AimaTestLogMapper, AimaT
@Override
public AimaTestLog selectAimaTestLogById(String id) {
return aimaTestLogMapper.selectAimaTestLogById(id);
return AimaLogExecutionContext.hydrate(aimaTestLogMapper.selectAimaTestLogById(id));
}
@Override
public List<AimaTestLog> selectAimaTestLogList(AimaTestLog aimaTestLog) {
return aimaTestLogMapper.selectAimaTestLogList(aimaTestLog);
List<AimaTestLog> logs = aimaTestLogMapper.selectAimaTestLogList(aimaTestLog);
logs.forEach(AimaLogExecutionContext::hydrate);
return logs;
}
@Override

View File

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

View File

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

View File

@ -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<AimaTestLog> 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;
}
}
}

View File

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

View File

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

View File

@ -64,6 +64,13 @@
<version>0.14.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.13.2</version>
<scope>test</scope>
</dependency>
</dependencies>
</project>

View File

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

View File

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

View File

@ -45,6 +45,11 @@ public interface EdgeAgvService {
*/
void clearFault(EdgeCommonVO edgeCommonVO);
/**
* 按指定速度平移固定距离。
*/
void translate(EdgeCommonVO edgeCommonVO, AgvUtils.AgvTranslation translation);
/**
* 导航到指定位置
*

View File

@ -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)
* 控制机械臂通过关节空间插值运动到目标关节位置

View File

@ -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<DeDeviceRegistration> deviceList(String robotId);

View File

@ -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) {

View File

@ -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,

View File

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

View File

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

View File

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

View File

@ -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", "视频智能分析"),
// 触控交互

View File

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

View File

@ -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);
// 从输入参数中获取坐标信息

View File

@ -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()) {

View File

@ -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
? "机械臂关节空间运动任务下发成功"
: "机械臂末端运动任务下发成功";

View File

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

View File

@ -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<JSONObject> 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