diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTaskInstanceController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTaskInstanceController.java index 5d29601..a77e4c2 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTaskInstanceController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTaskInstanceController.java @@ -14,6 +14,7 @@ import com.cmvr.common.enums.BusinessType; import com.cmvr.aima.domain.AimaTaskInstance; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; +import com.cmvr.aima.domain.vo.AimaTaskStartVO; import com.cmvr.aima.service.IAimaTaskInstanceService; import com.cmvr.common.utils.poi.ExcelUtil; import com.cmvr.common.core.page.TableDataInfo; @@ -82,8 +83,9 @@ public class AimaTaskInstanceController extends BaseController { @PreAuthorize("@ss.hasPermi('aima:taskinstance:edit')") @Log(title = "任务执行实例", businessType = BusinessType.UPDATE) @PutMapping("/start/{id}") - public AjaxResult start(@PathVariable("id") String id) { - return toAjax(aimaTaskInstanceService.startInstance(id)); + public AjaxResult start(@PathVariable("id") String id, + @RequestBody(required = false) AimaTaskStartVO startVO) { + return toAjax(aimaTaskInstanceService.startInstance(id, startVO)); } /** diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/inspection/InspectionTaskInstanceController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/inspection/InspectionTaskInstanceController.java index 549be6a..c00405e 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/inspection/InspectionTaskInstanceController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/inspection/InspectionTaskInstanceController.java @@ -21,6 +21,7 @@ import com.cmvr.common.core.controller.BaseController; import com.cmvr.common.core.domain.AjaxResult; import com.cmvr.common.enums.BusinessType; import com.cmvr.inspection.domain.InspectionTaskInstance; +import com.cmvr.inspection.domain.vo.InspectionTaskStartVO; import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo; import com.cmvr.inspection.service.IInspectionTaskInstanceService; import com.cmvr.common.utils.poi.ExcelUtil; @@ -121,9 +122,10 @@ public class InspectionTaskInstanceController extends BaseController @PreAuthorize("@ss.hasPermi('inspection:taskInstance:edit')") @Log(title = "巡检任务执行实例", businessType = BusinessType.UPDATE) @PutMapping("/start/{id}") - public AjaxResult start(@PathVariable("id") String id) + public AjaxResult start(@PathVariable("id") String id, + @RequestBody(required = false) InspectionTaskStartVO startVO) { - return toAjax(inspectionTaskInstanceService.startInstance(id)); + return toAjax(inspectionTaskInstanceService.startInstance(id, startVO)); } /** diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/domain/vo/AimaTaskStartVO.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/domain/vo/AimaTaskStartVO.java new file mode 100644 index 0000000..46dccaf --- /dev/null +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/domain/vo/AimaTaskStartVO.java @@ -0,0 +1,20 @@ +package com.cmvr.aima.domain.vo; + +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import lombok.Data; + +import java.util.ArrayList; +import java.util.List; + +@Data +@ApiModel("爱玛任务启动参数") +public class AimaTaskStartVO { + + @ApiModelProperty("兼容单机器人执行ID") + private String robotId; + + @ApiModelProperty("机器人角色运行分配") + private List runtimeAssignments = new ArrayList<>(); +} diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/IAimaTaskInstanceService.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/IAimaTaskInstanceService.java index 13d6bc8..cc17de5 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/IAimaTaskInstanceService.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/IAimaTaskInstanceService.java @@ -6,6 +6,7 @@ import com.baomidou.mybatisplus.spring.service.IService; import com.cmvr.aima.domain.AimaTaskInstance; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; +import com.cmvr.aima.domain.vo.AimaTaskStartVO; /** * 爱玛任务执行实例Service接口 @@ -28,7 +29,7 @@ public interface IAimaTaskInstanceService extends IService { * @param id 任务实例ID * @return 结果 */ - int startInstance(String id); + int startInstance(String id, AimaTaskStartVO startVO); /** * 暂停任务实例 diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java index 7021296..b5fef16 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java @@ -6,6 +6,7 @@ import com.alibaba.fastjson2.JSONObject; import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; +import com.cmvr.aima.domain.vo.AimaTaskStartVO; import com.cmvr.common.exception.ServiceException; import com.cmvr.common.utils.SecurityUtils; import com.cmvr.test.flow.control.FlowControlService; @@ -119,7 +120,7 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl runtimeAssignments = new ArrayList<>(); +} diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/IInspectionTaskInstanceService.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/IInspectionTaskInstanceService.java index 39ceba6..0b14159 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/IInspectionTaskInstanceService.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/IInspectionTaskInstanceService.java @@ -4,6 +4,7 @@ import java.util.List; import com.baomidou.mybatisplus.spring.service.IService; import com.cmvr.inspection.domain.InspectionTaskInstance; +import com.cmvr.inspection.domain.vo.InspectionTaskStartVO; import com.cmvr.inspection.domain.dto.InspectionTaskInstanceQuery; import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo; /** @@ -61,7 +62,7 @@ public interface IInspectionTaskInstanceService extends IService robotIds = new ArrayList<>(); + + /** 全部机器人锁键,用于多机器人任务的原子占用和释放。 */ + private List lockKeys = new ArrayList<>(); + + /** 任务启动时冻结的机器人角色分配。 */ + private List runtimeAssignments = new ArrayList<>(); + /** * 当前运行状态(运行中/已完成/失败/暂停/终止) */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContextManager.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContextManager.java index d4a9f54..8d357af 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContextManager.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskContextManager.java @@ -1,6 +1,7 @@ package com.cmvr.test.flow.context; import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; import com.cmvr.test.enums.TaskStatusEnum; import com.cmvr.test.enums.TerminalStatusEnum; import lombok.Getter; @@ -8,6 +9,8 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.util.List; +import java.util.ArrayList; +import java.util.Collections; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -32,30 +35,41 @@ public class TaskContextManager { /** * 注册任务上下文(任务启动时调用) */ - public void register(TaskContext context) { - instContextMap.put(context.getInstId(), context); - if (context.getLockKey() != null) { - robotInstMap.put(context.getLockKey(), context.getInstId()); + public synchronized void register(TaskContext context) { + List lockKeys = effectiveLockKeys(context); + for (String lockKey : lockKeys) { + if (robotInstMap.containsKey(lockKey)) { + throw new GlobalException("当前任务中的机器人正在执行其他任务,请稍后再试"); + } } + instContextMap.put(context.getInstId(), context); + lockKeys.forEach(lockKey -> robotInstMap.put(lockKey, context.getInstId())); context.setTerminalStatus(TerminalStatusEnum.RUNNING); - log.info("注册任务上下文:instId={}, robotId={}", context.getInstId(), context.getRobotId()); + log.info("注册任务上下文:instId={}, robotIds={}", context.getInstId(), context.getRobotIds()); } /** * 注销任务上下文(任务结束 / 终止后调用) */ - public void unregister(String instId) { + public synchronized void unregister(String instId) { TaskContext ctx = instContextMap.remove(instId); if (ctx != null) { ctx.setTerminalStatus(TerminalStatusEnum.READY); ctx.setTerminalLocked(false); - if (ctx.getLockKey() != null) { - robotInstMap.remove(ctx.getLockKey()); - } - log.info("注销任务上下文:instId={}, robotId={}", instId, ctx.getRobotId()); + effectiveLockKeys(ctx).forEach(lockKey -> robotInstMap.remove(lockKey, instId)); + log.info("注销任务上下文:instId={}, robotIds={}", instId, ctx.getRobotIds()); } } + private List effectiveLockKeys(TaskContext context) { + if (context.getLockKeys() != null && !context.getLockKeys().isEmpty()) { + return new ArrayList<>(context.getLockKeys()); + } + return context.getLockKey() == null + ? Collections.emptyList() + : Collections.singletonList(context.getLockKey()); + } + /** * 获取任务上下文 */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java index e555f30..8e1dac0 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java @@ -22,6 +22,7 @@ import org.springframework.stereotype.Component; import jakarta.annotation.Resource; import java.util.ArrayList; import java.util.List; +import java.util.Collections; import java.util.concurrent.CountDownLatch; import java.util.stream.Collectors; @@ -55,7 +56,7 @@ public class FlowControlService { throw new GlobalException("任务已终止,无需重复操作"); } // 异步终止 - executor.execute(() -> edgeSystemService.stopAll(ctx.getRobotId())); + stopAllRobots(ctx); // 统一记录日志 + 设置上下文状态 + 数据库状态 taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED); @@ -80,7 +81,7 @@ public class FlowControlService { throw new GlobalException("任务已处于暂停状态"); } // 异步停止终端 - executor.execute(() -> edgeSystemService.stopAll(ctx.getRobotId())); + stopAllRobots(ctx); ctx.setPaused(true); ctx.setTerminalStatus(TerminalStatusEnum.RUNNING); ctx.setStatus(TaskStatusEnum.PAUSED); @@ -196,4 +197,14 @@ public class FlowControlService { // 发布任务恢复事件 taskInstHolder.syncStatusAsResumed(instId); } + + private void stopAllRobots(TaskContext context) { + List robotIds = context.getRobotIds() == null || context.getRobotIds().isEmpty() + ? Collections.singletonList(context.getRobotId()) + : context.getRobotIds(); + robotIds.stream() + .filter(robotId -> robotId != null && !robotId.isBlank()) + .distinct() + .forEach(robotId -> executor.execute(() -> edgeSystemService.stopAll(robotId))); + } } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java index 261ffa1..249621f 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java @@ -40,6 +40,7 @@ public class FlowItemExecutor { private final TaskInstHolder taskInstHolder; private final FlowExecutionChainBuilder chainBuilder; private final TaskThreadRegistry taskThreadRegistry; + private final FlowRuntimeAssignmentResolver runtimeAssignmentResolver; /** * 执行主流程(每个检测项调用一次) @@ -226,7 +227,7 @@ public class FlowItemExecutor { } @NotNull - private static TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node, + private TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node, TaskNodeExecuteMessage rootMessage, String nodeId, String nodeName, List iterations) { TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); @@ -245,6 +246,7 @@ public class FlowItemExecutor { message.setLoopNum(rootMessage.getLoopNum()); message.setLoopArray(rootMessage.getLoopArray()); + runtimeAssignmentResolver.resolveNode(node, message); return message; } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolver.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolver.java new file mode 100644 index 0000000..72dcfe4 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolver.java @@ -0,0 +1,88 @@ +package com.cmvr.test.flow.runtime.engine; + +import cn.hutool.core.util.StrUtil; +import com.alibaba.fastjson2.JSONArray; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; +import com.cmvr.test.model.vo.RobotRuntimeMemberVO; +import org.springframework.stereotype.Component; + +import java.util.List; + +/** Resolves a logical role/device use to the frozen physical assignment of this task. */ +@Component +public class FlowRuntimeAssignmentResolver { + + public void resolveNode(FlowNodeWrapper node, TaskNodeExecuteMessage message) { + JSONObject executionTarget = node.getRawProperties().getJSONObject("executionTarget"); + if (executionTarget == null) { + return; + } + + String roleKey = executionTarget.getString("roleKey"); + if (StrUtil.isBlank(roleKey)) { + throw new GlobalException("节点[" + node.getNodeName() + "]未配置机器人角色"); + } + // 兼容旧调用方:只有robotId时继续使用根消息中的机器人和节点原有deviceId。 + if (message.getRuntimeAssignments() == null || message.getRuntimeAssignments().isEmpty()) { + message.setRoleKey(roleKey); + return; + } + ResourceRuntimeAssignmentVO assignment = findAssignment(message.getRuntimeAssignments(), roleKey); + if (assignment == null || assignment.getMembers() == null || assignment.getMembers().isEmpty()) { + throw new GlobalException("机器人角色[" + roleKey + "]未安排执行机器人"); + } + + // 第一阶段实现一个角色绑定一台机器人;members结构为后续角色池调度保留。 + RobotRuntimeMemberVO member = assignment.getMembers().get(0); + if (StrUtil.isBlank(member.getRobotId())) { + throw new GlobalException("机器人角色[" + roleKey + "]的机器人ID不能为空"); + } + message.setRoleKey(roleKey); + message.setRoleMemberKey(member.getMemberKey()); + message.setRobotId(member.getRobotId()); + + JSONObject resolvedDevices = new JSONObject(); + resolvedDevices.put("robotId", member.getRobotId()); + String deviceSlotKey = executionTarget.getString("deviceSlotKey"); + if (StrUtil.isNotBlank(deviceSlotKey)) { + String deviceId = member.getDevices() == null ? null : member.getDevices().get(deviceSlotKey); + if (StrUtil.isBlank(deviceId)) { + throw new GlobalException("机器人角色[" + roleKey + "]未分配设备用途[" + deviceSlotKey + "]"); + } + resolvedDevices.put("deviceId", deviceId); + } + + // 兼容Schema V2早期的parameterName/slotKey结构。 + JSONArray deviceBindings = executionTarget.getJSONArray("deviceBindings"); + if (deviceBindings != null) { + for (int i = 0; i < deviceBindings.size(); i++) { + JSONObject binding = deviceBindings.getJSONObject(i); + String parameterName = binding.getString("parameterName"); + String slotKey = binding.getString("slotKey"); + String deviceId = member.getDevices() == null ? null : member.getDevices().get(slotKey); + if (StrUtil.isBlank(parameterName) || StrUtil.isBlank(slotKey)) { + throw new GlobalException("节点[" + node.getNodeName() + "]的设备用途配置不完整"); + } + if (StrUtil.isBlank(deviceId)) { + throw new GlobalException("机器人角色[" + roleKey + "]未分配设备用途[" + slotKey + "]"); + } + resolvedDevices.put(parameterName, deviceId); + } + } + message.setResolvedDeviceBindings(resolvedDevices); + } + + private ResourceRuntimeAssignmentVO findAssignment(List assignments, String roleKey) { + if (assignments == null) { + return null; + } + return assignments.stream() + .filter(item -> roleKey.equals(item.getRoleKey())) + .findFirst() + .orElse(null); + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeDefinitionValidator.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeDefinitionValidator.java new file mode 100644 index 0000000..150ddd8 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeDefinitionValidator.java @@ -0,0 +1,130 @@ +package com.cmvr.test.flow.runtime.engine; + +import cn.hutool.core.util.StrUtil; +import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONArray; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; +import com.cmvr.test.model.vo.RobotRuntimeMemberVO; +import lombok.RequiredArgsConstructor; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.stereotype.Component; + +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; + +/** Validates Schema V2 role references before a task is submitted asynchronously. */ +@Component +@RequiredArgsConstructor +public class FlowRuntimeDefinitionValidator { + + private final ObjectProvider deviceValidatorProvider; + + public void validate(String flowData, List assignments) { + JSONObject flow = JSON.parseObject(flowData); + JSONArray nodes = flow.getJSONArray("nodes"); + if (nodes == null) { + throw new GlobalException("流程节点数据不能为空"); + } + JSONObject startProperties = null; + for (int i = 0; i < nodes.size(); i++) { + JSONObject node = nodes.getJSONObject(i); + if ("start".equalsIgnoreCase(node.getString("type"))) { + startProperties = node.getJSONObject("properties"); + break; + } + } + JSONArray roles = startProperties == null ? null : startProperties.getJSONArray("resourceRoles"); + if (roles == null || roles.isEmpty()) { + return; // 旧流程继续使用单robotId协议。 + } + + Map assignmentByRole = new HashMap<>(); + if (assignments != null) { + assignments.forEach(item -> assignmentByRole.put(item.getRoleKey(), item)); + } + Map> slotKeysByRole = new HashMap<>(); + + for (int i = 0; i < roles.size(); i++) { + JSONObject role = roles.getJSONObject(i); + String roleKey = role.getString("roleKey"); + if (StrUtil.isBlank(roleKey) || slotKeysByRole.containsKey(roleKey)) { + throw new GlobalException("流程机器人角色标识为空或重复"); + } + ResourceRuntimeAssignmentVO assignment = assignmentByRole.get(roleKey); + RobotRuntimeMemberVO member = firstMember(assignment); + boolean required = !Boolean.FALSE.equals(role.getBoolean("required")); + if (required && (member == null || StrUtil.isBlank(member.getRobotId()))) { + throw new GlobalException("机器人角色[" + role.getString("displayName") + "]未安排执行机器人"); + } + + Set slotKeys = new HashSet<>(); + JSONArray slots = role.getJSONArray("deviceSlots"); + if (slots != null) { + for (int j = 0; j < slots.size(); j++) { + JSONObject slot = slots.getJSONObject(j); + String slotKey = slot.getString("slotKey"); + if (StrUtil.isBlank(slotKey) || !slotKeys.add(slotKey)) { + throw new GlobalException("机器人角色[" + roleKey + "]的设备用途标识为空或重复"); + } + String deviceId = member == null || member.getDevices() == null + ? null : member.getDevices().get(slotKey); + boolean slotRequired = !Boolean.FALSE.equals(slot.getBoolean("required")); + boolean roleAssigned = member != null && StrUtil.isNotBlank(member.getRobotId()); + if (slotRequired && (required || roleAssigned) && StrUtil.isBlank(deviceId)) { + throw new GlobalException("设备用途[" + slot.getString("displayName") + "]未分配真实设备"); + } + validatePhysicalDevice(member, deviceId, slot.getInteger("deviceKind")); + } + } + slotKeysByRole.put(roleKey, slotKeys); + } + + for (int i = 0; i < nodes.size(); i++) { + JSONObject node = nodes.getJSONObject(i); + JSONObject properties = node.getJSONObject("properties"); + JSONObject target = properties == null ? null : properties.getJSONObject("executionTarget"); + if (target == null) { + continue; + } + String roleKey = target.getString("roleKey"); + Set slotKeys = slotKeysByRole.get(roleKey); + if (slotKeys == null) { + throw new GlobalException("节点[" + properties.getString("name") + "]引用了不存在的机器人角色"); + } + String deviceSlotKey = target.getString("deviceSlotKey"); + if (StrUtil.isNotBlank(deviceSlotKey) && !slotKeys.contains(deviceSlotKey)) { + throw new GlobalException("节点[" + properties.getString("name") + "]引用了不存在的设备用途"); + } + // 兼容Schema V2早期的deviceBindings结构。 + JSONArray bindings = target.getJSONArray("deviceBindings"); + if (bindings != null) { + for (int j = 0; j < bindings.size(); j++) { + String slotKey = bindings.getJSONObject(j).getString("slotKey"); + if (!slotKeys.contains(slotKey)) { + throw new GlobalException("节点[" + properties.getString("name") + "]引用了不存在的设备用途"); + } + } + } + } + } + + private RobotRuntimeMemberVO firstMember(ResourceRuntimeAssignmentVO assignment) { + return assignment == null || assignment.getMembers() == null || assignment.getMembers().isEmpty() + ? null : assignment.getMembers().get(0); + } + + private void validatePhysicalDevice(RobotRuntimeMemberVO member, String deviceId, Integer expectedKind) { + if (member == null || StrUtil.isBlank(deviceId)) { + return; + } + RobotDeviceAssignmentValidator validator = deviceValidatorProvider.getIfAvailable(); + if (validator != null) { + validator.validate(member.getRobotId(), deviceId, expectedKind); + } + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java index 5542f49..37a74bb 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java @@ -48,6 +48,7 @@ public class FlowTaskEngine { message.setInstId(instId); message.setTaskId(taskId); message.setRobotId(robotId); + message.setRuntimeAssignments(context.getRuntimeAssignments()); message.setItemId(itemId); if (runMode.equals(RunModeEnum.VI_PROJECT)) { message.setLoopDetail(item.getSchemeInfo()); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java index 5ec1634..851217f 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java @@ -18,6 +18,8 @@ import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO; import com.cmvr.test.model.vo.TeTaskExecuteNormalVO; import com.cmvr.test.model.vo.TeTaskExecuteProjectVO; import com.cmvr.test.model.vo.TeTaskExecuteTrailVO; +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; +import com.cmvr.test.model.vo.RobotRuntimeMemberVO; import com.cmvr.test.service.ITeDetectionItemService; import com.cmvr.test.service.ITeTaskInstService; import com.cmvr.test.service.ITeTaskOrchestrationService; @@ -28,6 +30,11 @@ import org.springframework.beans.factory.ObjectProvider; import org.springframework.transaction.annotation.Transactional; import java.util.List; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.Map; +import java.util.Set; /** * 任务执行入口 @@ -43,14 +50,16 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { private final TaskInstHolder taskInstHolder; private final FlowTaskAsyncDispatcher flowTaskAsyncDispatcher; private final ObjectProvider targetResolverProvider; + private final FlowRuntimeDefinitionValidator runtimeDefinitionValidator; @Transactional(rollbackFor = Exception.class) public String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO) { - RobotExecutionTarget target = resolveTarget(taskExecuteNormalVO.getRobotId()); + ResolvedRuntimeResources resources = resolveResources( + taskExecuteNormalVO.getRobotId(), taskExecuteNormalVO.getRuntimeAssignments()); return executeTaskInternal( - target, + resources, taskExecuteNormalVO.getTaskId(), taskExecuteNormalVO.getModuleCode(), RunModeEnum.NORMAL, @@ -62,18 +71,15 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { @Transactional(rollbackFor = Exception.class) public String executeTrialTask(TeTaskExecuteTrailVO taskExecuteTrailVO) { JSONObject runParams = taskExecuteTrailVO.getRunParams(); - RobotExecutionTarget target = resolveTarget(taskExecuteTrailVO.getRobotId()); + ResolvedRuntimeResources resources = resolveResources( + taskExecuteTrailVO.getRobotId(), taskExecuteTrailVO.getRuntimeAssignments()); + RobotExecutionTarget target = resources.primaryTarget(); String robotId = target.getRobotId(); - String lockKey = target.lockKey(); - if (lockKey == null) { - throw new GlobalException("机器人ID不能为空"); - } - if (taskInstHolder.isRobotLocked(lockKey)) { - throw new GlobalException("当前机器人正在执行其他任务,请稍后再试"); - } + assertResourcesAvailable(resources); // 构建试运行任务定义(临时任务) String itemId = taskExecuteTrailVO.getItemId(); + runtimeDefinitionValidator.validate(taskExecuteTrailVO.getFlowData(), resources.assignments()); TeDetectItemDeployFlowVO deployFlowVO = new TeDetectItemDeployFlowVO(); deployFlowVO.setId(itemId); @@ -96,7 +102,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { List details = CollUtil.newArrayList(teQueryTaskDetailDTO); // 注册上下文 - registerTaskContext(instId, taskId, null, target, itemId, RunModeEnum.TRIAL, runParams); + registerTaskContext(instId, taskId, null, resources, itemId, RunModeEnum.TRIAL, runParams); // 异步调度执行(与正式任务一致,只是模式为 TRIAL) flowTaskAsyncDispatcher.submit(instId, taskId, robotId, RunModeEnum.TRIAL, details); @@ -116,9 +122,10 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { @Transactional(rollbackFor = Exception.class) public String executeProjectTask(TeTaskExecuteProjectVO taskExecuteProjectVO) { - RobotExecutionTarget target = resolveTarget(taskExecuteProjectVO.getRobotId()); + ResolvedRuntimeResources resources = resolveResources( + taskExecuteProjectVO.getRobotId(), taskExecuteProjectVO.getRuntimeAssignments()); return executeTaskInternal( - target, + resources, taskExecuteProjectVO.getProjectId(), null, RunModeEnum.VI_PROJECT, @@ -130,19 +137,18 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { @Transactional(rollbackFor = Exception.class) public String executeTaskInternal(String robotId, String taskId, String moduleCode, RunModeEnum runMode, JSONObject runParams, Object originalVO) { - return executeTaskInternal(new RobotExecutionTarget(robotId), taskId, + return executeTaskInternal(resolveResources(robotId, null), taskId, moduleCode, runMode, runParams, originalVO); } - private String executeTaskInternal(RobotExecutionTarget target, String taskId, String moduleCode, + private String executeTaskInternal(ResolvedRuntimeResources resources, String taskId, String moduleCode, RunModeEnum runMode, JSONObject runParams, Object originalVO) { + RobotExecutionTarget target = resources.primaryTarget(); String robotId = target.getRobotId(); - String lockKey = target.lockKey(); - if (lockKey != null && taskInstHolder.isRobotLocked(lockKey)) { - throw new GlobalException("当前机器人正在执行其他任务,请稍后再试"); - } + assertResourcesAvailable(resources); // 查询任务详情(检测项信息) List details = taskOrchestrationService.queryTaskDetail(taskId); + details.forEach(detail -> runtimeDefinitionValidator.validate(detail.getFlowData(), resources.assignments())); int count = (int) details.stream().filter(dto -> ObjUtil.equals(dto.getIsDeploy(), "1")).count(); if (count > 0) { throw new GlobalException("当前任务存在未发布的检测项!"); @@ -153,7 +159,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { try { // 注册上下文 - registerTaskContext(instId, taskId, moduleCode, target, null, runMode, runParams); + registerTaskContext(instId, taskId, moduleCode, resources, null, runMode, runParams); // 异步提交执行 flowTaskAsyncDispatcher.submit(instId, taskId, robotId, runMode, details); @@ -188,8 +194,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { /** * 注册上下文 */ - private void registerTaskContext(String instId, String taskId, String moduleCode, RobotExecutionTarget target, + private void registerTaskContext(String instId, String taskId, String moduleCode, ResolvedRuntimeResources resources, String itemId, RunModeEnum mode, JSONObject runParams) { + RobotExecutionTarget target = resources.primaryTarget(); TaskContext ctx = new TaskContext(); ctx.setInstId(instId); ctx.setTaskId(taskId); @@ -197,6 +204,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { ctx.setRunMode(mode); ctx.setRobotId(target.getRobotId()); ctx.setLockKey(target.lockKey()); + ctx.setRobotIds(resources.targets().stream().map(RobotExecutionTarget::getRobotId).filter(StrUtil::isNotBlank).toList()); + ctx.setLockKeys(resources.targets().stream().map(RobotExecutionTarget::lockKey).filter(ObjUtil::isNotEmpty).toList()); + ctx.setRuntimeAssignments(resources.assignments()); ctx.setItemId(itemId); ctx.setStatus(TaskStatusEnum.RUNNING); ctx.setTerminalLocked(true); @@ -216,4 +226,58 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { // 纯逻辑工作流不调用边缘设备,可不绑定机器人。 return new RobotExecutionTarget(null); } + + private ResolvedRuntimeResources resolveResources(String legacyRobotId, + List assignments) { + List normalizedAssignments = assignments == null + ? new ArrayList<>() : assignments; + Map targetsByRobot = new LinkedHashMap<>(); + Set roleKeys = new LinkedHashSet<>(); + + for (ResourceRuntimeAssignmentVO assignment : normalizedAssignments) { + if (assignment == null || StrUtil.isBlank(assignment.getRoleKey())) { + throw new GlobalException("机器人角色标识不能为空"); + } + if (!roleKeys.add(assignment.getRoleKey())) { + throw new GlobalException("机器人角色重复:" + assignment.getRoleKey()); + } + if (assignment.getMembers() == null) { + continue; + } + for (RobotRuntimeMemberVO member : assignment.getMembers()) { + if (member == null || StrUtil.isBlank(member.getRobotId())) { + continue; + } + RobotExecutionTarget target = resolveTarget(member.getRobotId()); + member.setRobotId(target.getRobotId()); + targetsByRobot.putIfAbsent(target.getRobotId(), target); + } + } + + if (targetsByRobot.isEmpty() && StrUtil.isNotBlank(legacyRobotId)) { + RobotExecutionTarget target = resolveTarget(legacyRobotId); + targetsByRobot.put(target.getRobotId(), target); + } + List targets = new ArrayList<>(targetsByRobot.values()); + if (targets.isEmpty()) { + targets.add(new RobotExecutionTarget(null)); + } + return new ResolvedRuntimeResources(targets, normalizedAssignments); + } + + private void assertResourcesAvailable(ResolvedRuntimeResources resources) { + for (RobotExecutionTarget target : resources.targets()) { + String lockKey = target.lockKey(); + if (lockKey != null && taskInstHolder.isRobotLocked(lockKey)) { + throw new GlobalException("当前任务中的机器人正在执行其他任务,请稍后再试"); + } + } + } + + private record ResolvedRuntimeResources(List targets, + List assignments) { + private RobotExecutionTarget primaryTarget() { + return targets.get(0); + } + } } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/RobotDeviceAssignmentValidator.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/RobotDeviceAssignmentValidator.java new file mode 100644 index 0000000..f4fdf2f --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/RobotDeviceAssignmentValidator.java @@ -0,0 +1,7 @@ +package com.cmvr.test.flow.runtime.engine; + +/** Optional project adapter used to verify that a physical device belongs to a robot. */ +public interface RobotDeviceAssignmentValidator { + + void validate(String robotId, String deviceId, Integer expectedDeviceKind); +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowInputPrepareInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowInputPrepareInterceptor.java index 6c793c8..396a0b9 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowInputPrepareInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowInputPrepareInterceptor.java @@ -39,6 +39,16 @@ public class FlowInputPrepareInterceptor extends AbstractFlowMsgPreInterceptor { // 调用已有的 prepare 方法 JSONObject inputParams = paramPreparer.prepare(graph, node, message); + if (message.getResolvedDeviceBindings() != null) { + inputParams.putAll(message.getResolvedDeviceBindings()); + } + if (message.getRoleKey() != null) { + inputParams.put("executorRoleKey", message.getRoleKey()); + } + if (message.getRoleMemberKey() != null) { + inputParams.put("executorMemberKey", message.getRoleMemberKey()); + } + message.setInputParams(inputParams); log.debug("[InputPrepare] 节点 {} 准备输入参数: {}", message.getNodeId(), inputParams); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java index 8a341cd..009e1a9 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java @@ -3,6 +3,7 @@ package com.cmvr.test.flow.runtime.message; import com.alibaba.fastjson2.JSONObject; import com.cmvr.test.enums.ActionEnum; import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; import lombok.Data; import java.util.ArrayList; @@ -82,6 +83,18 @@ public class TaskNodeExecuteMessage { */ private String robotId; + /** 流程机器人角色标识。 */ + private String roleKey; + + /** 本次任务内被选中的角色成员。 */ + private String roleMemberKey; + + /** 冻结的本次执行角色分配。 */ + private List runtimeAssignments = new ArrayList<>(); + + /** 节点参数名到真实设备ID的解析结果。 */ + private JSONObject resolvedDeviceBindings = new JSONObject(); + /** * 节点入参 */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/ResourceRuntimeAssignmentVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/ResourceRuntimeAssignmentVO.java new file mode 100644 index 0000000..c61b924 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/ResourceRuntimeAssignmentVO.java @@ -0,0 +1,19 @@ +package com.cmvr.test.model.vo; + +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import lombok.Data; + +import java.util.ArrayList; +import java.util.List; + +@Data +@ApiModel("流程资源运行分配") +public class ResourceRuntimeAssignmentVO { + + @ApiModelProperty("流程内稳定的机器人角色标识") + private String roleKey; + + @ApiModelProperty("本次执行绑定的机器人成员") + private List members = new ArrayList<>(); +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/RobotRuntimeMemberVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/RobotRuntimeMemberVO.java new file mode 100644 index 0000000..0258ca7 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/RobotRuntimeMemberVO.java @@ -0,0 +1,22 @@ +package com.cmvr.test.model.vo; + +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import lombok.Data; + +import java.util.LinkedHashMap; +import java.util.Map; + +@Data +@ApiModel("机器人角色运行成员") +public class RobotRuntimeMemberVO { + + @ApiModelProperty("本次任务内的成员标识") + private String memberKey; + + @ApiModelProperty("机器人永久唯一ID") + private String robotId; + + @ApiModelProperty("设备用途标识到真实设备ID的映射") + private Map devices = new LinkedHashMap<>(); +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteNormalVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteNormalVO.java index 83c9bc2..92f914d 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteNormalVO.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteNormalVO.java @@ -10,6 +10,9 @@ import lombok.NoArgsConstructor; import jakarta.validation.constraints.NotEmpty; +import java.util.ArrayList; +import java.util.List; + @Data @ApiModel("正常任务执行VO") @Builder @@ -26,6 +29,11 @@ public class TeTaskExecuteNormalVO { @ApiModelProperty("机器人永久唯一ID") private String robotId; + @ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId") + @Builder.Default + private List runtimeAssignments = new ArrayList<>(); + @ApiModelProperty(value = "运行时参数") + @Builder.Default private JSONObject runParams = new JSONObject(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteProjectVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteProjectVO.java index 28d33ed..f896c90 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteProjectVO.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteProjectVO.java @@ -7,6 +7,9 @@ import lombok.Data; import jakarta.validation.constraints.NotEmpty; +import java.util.ArrayList; +import java.util.List; + @Data @ApiModel("项目执行VO") public class TeTaskExecuteProjectVO { @@ -18,6 +21,9 @@ public class TeTaskExecuteProjectVO { @ApiModelProperty("机器人永久唯一ID") private String robotId; + @ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId") + private List runtimeAssignments = new ArrayList<>(); + @ApiModelProperty(value = "运行时参数") private JSONObject runParams = new JSONObject(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteTrailVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteTrailVO.java index 796bff1..ff4aad0 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteTrailVO.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeTaskExecuteTrailVO.java @@ -7,6 +7,9 @@ import lombok.Data; import jakarta.validation.constraints.NotEmpty; +import java.util.ArrayList; +import java.util.List; + @Data @ApiModel("试运行任务执行VO") public class TeTaskExecuteTrailVO { @@ -20,9 +23,11 @@ public class TeTaskExecuteTrailVO { private String flowData; @ApiModelProperty("机器人永久唯一ID") - @NotEmpty(message = "机器人ID不能为空") private String robotId; + @ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId") + private List runtimeAssignments = new ArrayList<>(); + @ApiModelProperty(value = "运行时参数") private JSONObject runParams = new JSONObject(); } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolverTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolverTest.java new file mode 100644 index 0000000..dd980ea --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowRuntimeAssignmentResolverTest.java @@ -0,0 +1,68 @@ +package com.cmvr.test.flow.runtime.engine; + +import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; +import com.cmvr.test.model.vo.RobotRuntimeMemberVO; +import org.junit.Test; + +import java.util.List; +import java.util.Map; + +import static org.junit.Assert.assertEquals; + +public class FlowRuntimeAssignmentResolverTest { + + private final FlowRuntimeAssignmentResolver resolver = new FlowRuntimeAssignmentResolver(); + + @Test + public void resolvesRobotAndDeviceByRoleAndSlot() { + FlowNodeWrapper node = new FlowNodeWrapper(); + node.setNodeName("拍摄"); + node.setRawProperties(JSON.parseObject(""" + { + "executionTarget": { + "roleKey": "inspection", + "deviceSlotKey": "front_camera" + } + } + """)); + + RobotRuntimeMemberVO member = new RobotRuntimeMemberVO(); + member.setMemberKey("inspection_1"); + member.setRobotId("robot-2"); + member.setDevices(Map.of("front_camera", "camera-6")); + ResourceRuntimeAssignmentVO assignment = new ResourceRuntimeAssignmentVO(); + assignment.setRoleKey("inspection"); + assignment.setMembers(List.of(member)); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setRobotId("legacy-robot"); + message.setRuntimeAssignments(List.of(assignment)); + + resolver.resolveNode(node, message); + + assertEquals("inspection", message.getRoleKey()); + assertEquals("inspection_1", message.getRoleMemberKey()); + assertEquals("robot-2", message.getRobotId()); + assertEquals("camera-6", message.getResolvedDeviceBindings().getString("deviceId")); + } + + @Test + public void keepsLegacyRobotWhenAssignmentsAreAbsent() { + FlowNodeWrapper node = new FlowNodeWrapper(); + node.setNodeName("旧节点"); + node.setRawProperties(new JSONObject(Map.of( + "executionTarget", new JSONObject(Map.of("roleKey", "default_executor")) + ))); + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setRobotId("legacy-robot"); + + resolver.resolveNode(node, message); + + assertEquals("default_executor", message.getRoleKey()); + assertEquals("legacy-robot", message.getRobotId()); + } +} diff --git a/sql/inspection_module.sql b/sql/inspection_module.sql index c5b8e5e..6e6cdb0 100644 --- a/sql/inspection_module.sql +++ b/sql/inspection_module.sql @@ -61,7 +61,7 @@ CREATE TABLE `inspection_robot_device` ( `remark` varchar(500) DEFAULT NULL, PRIMARY KEY (`id`), KEY `idx_robot_device_robot_id` (`robot_id`), - KEY `idx_robot_device_identity` (`robot_id`, `device_id`), + UNIQUE KEY `uk_robot_device_identity` (`robot_id`, `device_id`), KEY `idx_robot_device_online` (`robot_id`, `online_status`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='机器人自动发现设备'; diff --git a/sql/robot_quic_identity_migration.sql b/sql/robot_quic_identity_migration.sql index b763003..cea8d44 100644 --- a/sql/robot_quic_identity_migration.sql +++ b/sql/robot_quic_identity_migration.sql @@ -35,7 +35,7 @@ create table inspection_robot_device ( remark varchar(500) null, primary key (id), index idx_robot_device_robot_id (robot_id), - index idx_robot_device_identity (robot_id, device_id), + unique index uk_robot_device_identity (robot_id, device_id), index idx_robot_device_online (robot_id, online_status) ) comment='Devices discovered from robot QUIC heartbeat'; diff --git a/sql/workflow_resource_role_v2.sql b/sql/workflow_resource_role_v2.sql new file mode 100644 index 0000000..dc4901d --- /dev/null +++ b/sql/workflow_resource_role_v2.sql @@ -0,0 +1,10 @@ +-- Workflow resource role V2 prerequisite. +-- Resolve any rows returned by this query before applying the unique constraint. +select robot_id, device_id, count(*) as duplicate_count +from inspection_robot_device +group by robot_id, device_id +having count(*) > 1; + +alter table inspection_robot_device + drop index idx_robot_device_identity, + add unique key uk_robot_device_identity (robot_id, device_id);