Compare commits
2 Commits
c51823801d
...
34ff31366d
| Author | SHA1 | Date | |
|---|---|---|---|
| 34ff31366d | |||
| 1543a6fb9c |
@ -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));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@ -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));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@ -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<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
}
|
||||
@ -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<AimaTaskInstance> {
|
||||
* @param id 任务实例ID
|
||||
* @return 结果
|
||||
*/
|
||||
int startInstance(String id);
|
||||
int startInstance(String id, AimaTaskStartVO startVO);
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
|
||||
@ -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<AimaTaskInstanceMap
|
||||
* 开始执行任务实例
|
||||
*/
|
||||
@Override
|
||||
public int startInstance(String id) {
|
||||
public int startInstance(String id, AimaTaskStartVO startVO) {
|
||||
AimaTaskInstance instance = aimaTaskInstanceMapper.selectAimaTaskInstanceById(id);
|
||||
if (instance == null) {
|
||||
throw new ServiceException("任务实例不存在");
|
||||
@ -148,6 +149,8 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
|
||||
TeTaskExecuteNormalVO.builder()
|
||||
.taskId(aimaTask.getTaskConfigId())
|
||||
.moduleCode(FlowExecutionModuleCodes.AIMA)
|
||||
.robotId(startVO == null ? null : startVO.getRobotId())
|
||||
.runtimeAssignments(startVO == null ? null : startVO.getRuntimeAssignments())
|
||||
.runParams(runParams)
|
||||
.build()
|
||||
);
|
||||
|
||||
@ -90,6 +90,9 @@ public class EdgeAgvServiceImpl implements EdgeAgvService {
|
||||
edgeCommonVO.getRobotId(), AgvServiceGrpc.AgvServiceBlockingStub.class);
|
||||
AgvCommand.AgvNavigateToStationCommand.Request request = AgvCommand.AgvNavigateToStationCommand.Request.newBuilder()
|
||||
.setHeader(EdgeCommonUtil.buildRequest(edgeCommonVO.getDeviceId()))
|
||||
.setOptions(AgvUtils.AgvMotionOptions.newBuilder()
|
||||
.setAsynchronous(true)
|
||||
.build())
|
||||
.setStationId(stationId)
|
||||
.build();
|
||||
executeGrpcCall(() -> stub.navigateToStation(request));
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@ -639,6 +639,37 @@ public final class AgvServiceGrpc {
|
||||
return getStopMappingMethod;
|
||||
}
|
||||
|
||||
private static volatile io.grpc.MethodDescriptor<cmvr.api.AgvCommand.AgvTranslateCommand.Request,
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> getTranslateMethod;
|
||||
|
||||
@io.grpc.stub.annotations.RpcMethod(
|
||||
fullMethodName = SERVICE_NAME + '/' + "translate",
|
||||
requestType = cmvr.api.AgvCommand.AgvTranslateCommand.Request.class,
|
||||
responseType = cmvr.api.AgvCommand.AgvTranslateCommand.Feedback.class,
|
||||
methodType = io.grpc.MethodDescriptor.MethodType.UNARY)
|
||||
public static io.grpc.MethodDescriptor<cmvr.api.AgvCommand.AgvTranslateCommand.Request,
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> getTranslateMethod() {
|
||||
io.grpc.MethodDescriptor<cmvr.api.AgvCommand.AgvTranslateCommand.Request, cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> getTranslateMethod;
|
||||
if ((getTranslateMethod = AgvServiceGrpc.getTranslateMethod) == null) {
|
||||
synchronized (AgvServiceGrpc.class) {
|
||||
if ((getTranslateMethod = AgvServiceGrpc.getTranslateMethod) == null) {
|
||||
AgvServiceGrpc.getTranslateMethod = getTranslateMethod =
|
||||
io.grpc.MethodDescriptor.<cmvr.api.AgvCommand.AgvTranslateCommand.Request, cmvr.api.AgvCommand.AgvTranslateCommand.Feedback>newBuilder()
|
||||
.setType(io.grpc.MethodDescriptor.MethodType.UNARY)
|
||||
.setFullMethodName(generateFullMethodName(SERVICE_NAME, "translate"))
|
||||
.setSampledToLocalTracing(true)
|
||||
.setRequestMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Request.getDefaultInstance()))
|
||||
.setResponseMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Feedback.getDefaultInstance()))
|
||||
.setSchemaDescriptor(new AgvServiceMethodDescriptorSupplier("translate"))
|
||||
.build();
|
||||
}
|
||||
}
|
||||
}
|
||||
return getTranslateMethod;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new async stub that supports all call types for the service
|
||||
*/
|
||||
@ -893,6 +924,16 @@ public final class AgvServiceGrpc {
|
||||
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getStopMappingMethod(), responseObserver);
|
||||
}
|
||||
|
||||
/**
|
||||
* <pre>
|
||||
* 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
|
||||
* </pre>
|
||||
*/
|
||||
public void translate(cmvr.api.AgvCommand.AgvTranslateCommand.Request request,
|
||||
io.grpc.stub.StreamObserver<cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> responseObserver) {
|
||||
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getTranslateMethod(), responseObserver);
|
||||
}
|
||||
|
||||
@java.lang.Override public final io.grpc.ServerServiceDefinition bindService() {
|
||||
return io.grpc.ServerServiceDefinition.builder(getServiceDescriptor())
|
||||
.addMethod(
|
||||
@ -1035,6 +1076,13 @@ public final class AgvServiceGrpc {
|
||||
cmvr.api.Common.CommandHeader.Request,
|
||||
cmvr.api.Common.CommandHeader.Feedback>(
|
||||
this, METHODID_STOP_MAPPING)))
|
||||
.addMethod(
|
||||
getTranslateMethod(),
|
||||
io.grpc.stub.ServerCalls.asyncUnaryCall(
|
||||
new MethodHandlers<
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Request,
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Feedback>(
|
||||
this, METHODID_TRANSLATE)))
|
||||
.build();
|
||||
}
|
||||
}
|
||||
@ -1278,6 +1326,17 @@ public final class AgvServiceGrpc {
|
||||
io.grpc.stub.ClientCalls.asyncUnaryCall(
|
||||
getChannel().newCall(getStopMappingMethod(), getCallOptions()), request, responseObserver);
|
||||
}
|
||||
|
||||
/**
|
||||
* <pre>
|
||||
* 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
|
||||
* </pre>
|
||||
*/
|
||||
public void translate(cmvr.api.AgvCommand.AgvTranslateCommand.Request request,
|
||||
io.grpc.stub.StreamObserver<cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> responseObserver) {
|
||||
io.grpc.stub.ClientCalls.asyncUnaryCall(
|
||||
getChannel().newCall(getTranslateMethod(), getCallOptions()), request, responseObserver);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@ -1500,6 +1559,16 @@ public final class AgvServiceGrpc {
|
||||
return io.grpc.stub.ClientCalls.blockingUnaryCall(
|
||||
getChannel(), getStopMappingMethod(), getCallOptions(), request);
|
||||
}
|
||||
|
||||
/**
|
||||
* <pre>
|
||||
* 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
|
||||
* </pre>
|
||||
*/
|
||||
public cmvr.api.AgvCommand.AgvTranslateCommand.Feedback translate(cmvr.api.AgvCommand.AgvTranslateCommand.Request request) {
|
||||
return io.grpc.stub.ClientCalls.blockingUnaryCall(
|
||||
getChannel(), getTranslateMethod(), getCallOptions(), request);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@ -1729,6 +1798,17 @@ public final class AgvServiceGrpc {
|
||||
return io.grpc.stub.ClientCalls.futureUnaryCall(
|
||||
getChannel().newCall(getStopMappingMethod(), getCallOptions()), request);
|
||||
}
|
||||
|
||||
/**
|
||||
* <pre>
|
||||
* 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
|
||||
* </pre>
|
||||
*/
|
||||
public com.google.common.util.concurrent.ListenableFuture<cmvr.api.AgvCommand.AgvTranslateCommand.Feedback> translate(
|
||||
cmvr.api.AgvCommand.AgvTranslateCommand.Request request) {
|
||||
return io.grpc.stub.ClientCalls.futureUnaryCall(
|
||||
getChannel().newCall(getTranslateMethod(), getCallOptions()), request);
|
||||
}
|
||||
}
|
||||
|
||||
private static final int METHODID_GET_RUNTIME_STATE = 0;
|
||||
@ -1751,6 +1831,7 @@ public final class AgvServiceGrpc {
|
||||
private static final int METHODID_START_MAPPING = 17;
|
||||
private static final int METHODID_STREAM_MAP = 18;
|
||||
private static final int METHODID_STOP_MAPPING = 19;
|
||||
private static final int METHODID_TRANSLATE = 20;
|
||||
|
||||
private static final class MethodHandlers<Req, Resp> implements
|
||||
io.grpc.stub.ServerCalls.UnaryMethod<Req, Resp>,
|
||||
@ -1849,6 +1930,10 @@ public final class AgvServiceGrpc {
|
||||
serviceImpl.stopMapping((cmvr.api.Common.CommandHeader.Request) request,
|
||||
(io.grpc.stub.StreamObserver<cmvr.api.Common.CommandHeader.Feedback>) responseObserver);
|
||||
break;
|
||||
case METHODID_TRANSLATE:
|
||||
serviceImpl.translate((cmvr.api.AgvCommand.AgvTranslateCommand.Request) request,
|
||||
(io.grpc.stub.StreamObserver<cmvr.api.AgvCommand.AgvTranslateCommand.Feedback>) responseObserver);
|
||||
break;
|
||||
default:
|
||||
throw new AssertionError();
|
||||
}
|
||||
@ -1930,6 +2015,7 @@ public final class AgvServiceGrpc {
|
||||
.addMethod(getStartMappingMethod())
|
||||
.addMethod(getStreamMapMethod())
|
||||
.addMethod(getStopMappingMethod())
|
||||
.addMethod(getTranslateMethod())
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
||||
@ -25,7 +25,7 @@ public final class AgvServiceOuterClass {
|
||||
java.lang.String[] descriptorData = {
|
||||
"\n\032cmvr/api/agv_service.proto\022\010cmvr.api\032\025" +
|
||||
"cmvr/api/common.proto\032\032cmvr/api/agv_comm" +
|
||||
"and.proto2\320\016\n\nAgvService\022f\n\017getRuntimeSt" +
|
||||
"and.proto2\254\017\n\nAgvService\022f\n\017getRuntimeSt" +
|
||||
"ate\022(.cmvr.api.AgvRuntimeStateCommand.Re" +
|
||||
"quest\032).cmvr.api.AgvRuntimeStateCommand." +
|
||||
"Feedback\022r\n\023getNavigationStatus\022,.cmvr.a" +
|
||||
@ -72,7 +72,9 @@ public final class AgvServiceOuterClass {
|
||||
"\032&.cmvr.api.AgvMapStreamCommand.Feedback" +
|
||||
"0\001\022P\n\013stopMapping\022\037.cmvr.api.CommandHead" +
|
||||
"er.Request\032 .cmvr.api.CommandHeader.Feed" +
|
||||
"backb\006proto3"
|
||||
"back\022Z\n\ttranslate\022%.cmvr.api.AgvTranslat" +
|
||||
"eCommand.Request\032&.cmvr.api.AgvTranslate" +
|
||||
"Command.Feedbackb\006proto3"
|
||||
};
|
||||
descriptor = com.google.protobuf.Descriptors.FileDescriptor
|
||||
.internalBuildGeneratedFileFrom(descriptorData,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@ -95,6 +95,23 @@ message AgvFollowPathCommand {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// 按指定速度执行固定距离平移命令。
|
||||
message AgvTranslateCommand {
|
||||
message Request {
|
||||
// 通用请求头。header.device_id 指定目标 AGV 设备。
|
||||
CommandHeader.Request header = 1;
|
||||
// 固定距离平移参数。
|
||||
cmvr.msgs.AgvTranslation translation = 2;
|
||||
}
|
||||
|
||||
message Feedback {
|
||||
// 仅表示控制器是否接受命令,不表示平移已经完成。
|
||||
CommandHeader.Feedback header = 1;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// 下发底盘速度命令。
|
||||
message AgvSetVelocityCommand {
|
||||
// 请求体。
|
||||
|
||||
@ -69,4 +69,8 @@ service AgvService {
|
||||
|
||||
// 停止当前建图/扫图会话。
|
||||
rpc stopMapping(CommandHeader.Request) returns (CommandHeader.Feedback);
|
||||
|
||||
|
||||
// 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
|
||||
rpc translate(AgvTranslateCommand.Request) returns (AgvTranslateCommand.Feedback);
|
||||
}
|
||||
|
||||
@ -14,6 +14,29 @@ message AgvPose2d {
|
||||
double theta = 3;
|
||||
}
|
||||
|
||||
// 固定距离平移使用的距离参考模式。
|
||||
enum AgvTranslationMode {
|
||||
// 根据底盘里程计算运动距离。
|
||||
AGV_TRANSLATION_MODE_ODOMETRY = 0;
|
||||
// 根据定位结果计算运动距离。
|
||||
AGV_TRANSLATION_MODE_LOCALIZATION = 1;
|
||||
}
|
||||
|
||||
|
||||
|
||||
// AGV 车体坐标系下的固定距离平移参数。
|
||||
message AgvTranslation {
|
||||
// 平移距离的绝对值,单位:米,必须大于 0。
|
||||
double distance = 1;
|
||||
// 车体 X 方向速度,单位:米/秒;正为向前,负为向后。
|
||||
double vx = 2;
|
||||
// 车体 Y 方向速度,单位:米/秒;正为向左,负为向右。
|
||||
double vy = 3;
|
||||
// 距离参考模式;默认使用里程模式。
|
||||
AgvTranslationMode mode = 4;
|
||||
}
|
||||
|
||||
|
||||
// AGV 车体坐标系下的平面速度。
|
||||
message AgvVelocity {
|
||||
// 车体 X 方向线速度,单位:米/秒。
|
||||
@ -24,6 +47,9 @@ message AgvVelocity {
|
||||
double wz = 3;
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
// AGV 电池状态。
|
||||
message AgvBatteryState {
|
||||
// 电量比例,范围:[0, 1],例如 0.8 表示 80%。
|
||||
|
||||
@ -0,0 +1,20 @@
|
||||
package com.cmvr.inspection.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 InspectionTaskStartVO {
|
||||
|
||||
@ApiModelProperty("兼容单机器人执行ID")
|
||||
private String robotId;
|
||||
|
||||
@ApiModelProperty("机器人角色运行分配")
|
||||
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
}
|
||||
@ -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<InspectionTaskI
|
||||
* @param id 任务实例主键
|
||||
* @return 结果
|
||||
*/
|
||||
int startInstance(String id);
|
||||
int startInstance(String id, InspectionTaskStartVO startVO);
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
|
||||
@ -0,0 +1,26 @@
|
||||
package com.cmvr.inspection.service;
|
||||
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.inspection.mapper.InspectionRobotDeviceMapper;
|
||||
import com.cmvr.test.flow.runtime.engine.RobotDeviceAssignmentValidator;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
@Service
|
||||
@RequiredArgsConstructor
|
||||
public class InspectionRobotDeviceAssignmentValidator implements RobotDeviceAssignmentValidator {
|
||||
|
||||
private final InspectionRobotDeviceMapper deviceMapper;
|
||||
|
||||
@Override
|
||||
public void validate(String robotId, String deviceId, Integer expectedDeviceKind) {
|
||||
var devices = deviceMapper.selectByIdentity(robotId, deviceId);
|
||||
if (devices.isEmpty()) {
|
||||
throw new GlobalException("设备[" + deviceId + "]不属于机器人[" + robotId + "]");
|
||||
}
|
||||
Integer actualKind = devices.get(0).getDeviceKind();
|
||||
if (expectedDeviceKind != null && !expectedDeviceKind.equals(actualKind)) {
|
||||
throw new GlobalException("设备[" + deviceId + "]的类型与流程设备用途不匹配");
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -25,6 +25,7 @@ import org.springframework.stereotype.Service;
|
||||
import com.cmvr.common.utils.DateUtils;
|
||||
import com.cmvr.inspection.mapper.InspectionTaskInstanceMapper;
|
||||
import com.cmvr.inspection.domain.InspectionTaskInstance;
|
||||
import com.cmvr.inspection.domain.vo.InspectionTaskStartVO;
|
||||
import com.cmvr.inspection.service.IInspectionTaskInstanceService;
|
||||
|
||||
/**
|
||||
@ -143,7 +144,7 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
|
||||
* 开始执行任务实例
|
||||
*/
|
||||
@Override
|
||||
public int startInstance(String id)
|
||||
public int startInstance(String id, InspectionTaskStartVO startVO)
|
||||
{
|
||||
InspectionTaskInstance instance = inspectionTaskInstanceMapper.selectInspectionTaskInstanceById(id);
|
||||
if (instance == null) {
|
||||
@ -169,7 +170,8 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
|
||||
// 调用真正的执行
|
||||
String insId = flowTaskRuntimeService.executeTask(TeTaskExecuteNormalVO.builder()
|
||||
.taskId(inspectionTask.getTaskConfigId())
|
||||
.robotId(robot.getRobotId())
|
||||
.robotId(startVO != null && startVO.getRobotId() != null ? startVO.getRobotId() : robot.getRobotId())
|
||||
.runtimeAssignments(startVO == null ? null : startVO.getRuntimeAssignments())
|
||||
.moduleCode(FlowExecutionModuleCodes.INSPECTION)
|
||||
.runParams(runParams)
|
||||
.build());
|
||||
|
||||
@ -7,11 +7,13 @@ import com.cmvr.test.enums.RunModeEnum;
|
||||
import com.cmvr.test.enums.TaskStatusEnum;
|
||||
import com.cmvr.test.enums.TerminalStatusEnum;
|
||||
import com.cmvr.test.flow.builder.TaskKeyBuilder;
|
||||
import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO;
|
||||
import lombok.Data;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
@ -47,6 +49,15 @@ public class TaskContext {
|
||||
/** 稳定的机器人锁键。 */
|
||||
private String lockKey;
|
||||
|
||||
/** 本次任务实际使用的全部机器人ID。 */
|
||||
private List<String> robotIds = new ArrayList<>();
|
||||
|
||||
/** 全部机器人锁键,用于多机器人任务的原子占用和释放。 */
|
||||
private List<String> lockKeys = new ArrayList<>();
|
||||
|
||||
/** 任务启动时冻结的机器人角色分配。 */
|
||||
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
|
||||
/**
|
||||
* 当前运行状态(运行中/已完成/失败/暂停/终止)
|
||||
*/
|
||||
|
||||
@ -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<String> 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<String> 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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取任务上下文
|
||||
*/
|
||||
|
||||
@ -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<String> 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)));
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<Integer> 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;
|
||||
}
|
||||
|
||||
|
||||
@ -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<ResourceRuntimeAssignmentVO> assignments, String roleKey) {
|
||||
if (assignments == null) {
|
||||
return null;
|
||||
}
|
||||
return assignments.stream()
|
||||
.filter(item -> roleKey.equals(item.getRoleKey()))
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
}
|
||||
}
|
||||
@ -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<RobotDeviceAssignmentValidator> deviceValidatorProvider;
|
||||
|
||||
public void validate(String flowData, List<ResourceRuntimeAssignmentVO> 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<String, ResourceRuntimeAssignmentVO> assignmentByRole = new HashMap<>();
|
||||
if (assignments != null) {
|
||||
assignments.forEach(item -> assignmentByRole.put(item.getRoleKey(), item));
|
||||
}
|
||||
Map<String, Set<String>> 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<String> 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<String> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -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());
|
||||
|
||||
@ -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<RobotExecutionTargetResolver> 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<TeQueryTaskDetailDTO> 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<TeQueryTaskDetailDTO> 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<ResourceRuntimeAssignmentVO> assignments) {
|
||||
List<ResourceRuntimeAssignmentVO> normalizedAssignments = assignments == null
|
||||
? new ArrayList<>() : assignments;
|
||||
Map<String, RobotExecutionTarget> targetsByRobot = new LinkedHashMap<>();
|
||||
Set<String> 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<RobotExecutionTarget> 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<RobotExecutionTarget> targets,
|
||||
List<ResourceRuntimeAssignmentVO> assignments) {
|
||||
private RobotExecutionTarget primaryTarget() {
|
||||
return targets.get(0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@ -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);
|
||||
}
|
||||
@ -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);
|
||||
|
||||
|
||||
@ -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<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
|
||||
/** 节点参数名到真实设备ID的解析结果。 */
|
||||
private JSONObject resolvedDeviceBindings = new JSONObject();
|
||||
|
||||
/**
|
||||
* 节点入参
|
||||
*/
|
||||
|
||||
@ -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<RobotRuntimeMemberVO> members = new ArrayList<>();
|
||||
}
|
||||
@ -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<String, String> devices = new LinkedHashMap<>();
|
||||
}
|
||||
@ -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<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
|
||||
@ApiModelProperty(value = "运行时参数")
|
||||
@Builder.Default
|
||||
private JSONObject runParams = new JSONObject();
|
||||
}
|
||||
|
||||
@ -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<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
|
||||
@ApiModelProperty(value = "运行时参数")
|
||||
private JSONObject runParams = new JSONObject();
|
||||
}
|
||||
|
||||
@ -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<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
|
||||
|
||||
@ApiModelProperty(value = "运行时参数")
|
||||
private JSONObject runParams = new JSONObject();
|
||||
}
|
||||
|
||||
@ -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());
|
||||
}
|
||||
}
|
||||
@ -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='机器人自动发现设备';
|
||||
|
||||
|
||||
@ -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';
|
||||
|
||||
|
||||
10
sql/workflow_resource_role_v2.sql
Normal file
10
sql/workflow_resource_role_v2.sql
Normal file
@ -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);
|
||||
Loading…
Reference in New Issue
Block a user