Compare commits

...

2 Commits

Author SHA1 Message Date
34ff31366d feat(task): 支持多机器人角色分配的任务执行功能
- 在任务启动接口中增加 AimaTaskStartVO 和 InspectionTaskStartVO 参数支持
- 实现多机器人角色分配机制,支持运行时动态分配机器人和设备资源
- 添加 FlowRuntimeAssignmentResolver 组件处理节点执行目标解析
- 更新任务执行引擎以支持多机器人并发控制和资源锁定
- 修改数据库约束将机器人设备关联改为唯一索引防止重复绑定
- 增强流程定义验证器确保运行时分配配置的完整性
- 扩展任务上下文管理支持多机器人 ID 和锁键跟踪
- 优化试运行和项目任务执行流程集成新的资源分配逻辑
2026-08-11 16:54:01 +08:00
1543a6fb9c feat(agv): 添加AGV固定距离平移功能
- 定义AgvTranslateCommand消息结构,包含请求和反馈
- 添加AgvTranslationMode枚举,支持里程和定位两种距离参考模式
- 定义AgvTranslation消息,包含距离、速度和平移模式参数
- 在AgvService服务中添加translate RPC方法
- 生成Java gRPC客户端和服务端代码实现
- 更新protobuf描述符以包含新的平移命令定义
2026-08-11 15:16:15 +08:00
37 changed files with 3878 additions and 187 deletions

View File

@ -14,6 +14,7 @@ import com.cmvr.common.enums.BusinessType;
import com.cmvr.aima.domain.AimaTaskInstance; import com.cmvr.aima.domain.AimaTaskInstance;
import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery;
import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo;
import com.cmvr.aima.domain.vo.AimaTaskStartVO;
import com.cmvr.aima.service.IAimaTaskInstanceService; import com.cmvr.aima.service.IAimaTaskInstanceService;
import com.cmvr.common.utils.poi.ExcelUtil; import com.cmvr.common.utils.poi.ExcelUtil;
import com.cmvr.common.core.page.TableDataInfo; import com.cmvr.common.core.page.TableDataInfo;
@ -82,8 +83,9 @@ public class AimaTaskInstanceController extends BaseController {
@PreAuthorize("@ss.hasPermi('aima:taskinstance:edit')") @PreAuthorize("@ss.hasPermi('aima:taskinstance:edit')")
@Log(title = "任务执行实例", businessType = BusinessType.UPDATE) @Log(title = "任务执行实例", businessType = BusinessType.UPDATE)
@PutMapping("/start/{id}") @PutMapping("/start/{id}")
public AjaxResult start(@PathVariable("id") String id) { public AjaxResult start(@PathVariable("id") String id,
return toAjax(aimaTaskInstanceService.startInstance(id)); @RequestBody(required = false) AimaTaskStartVO startVO) {
return toAjax(aimaTaskInstanceService.startInstance(id, startVO));
} }
/** /**

View File

@ -21,6 +21,7 @@ import com.cmvr.common.core.controller.BaseController;
import com.cmvr.common.core.domain.AjaxResult; import com.cmvr.common.core.domain.AjaxResult;
import com.cmvr.common.enums.BusinessType; import com.cmvr.common.enums.BusinessType;
import com.cmvr.inspection.domain.InspectionTaskInstance; import com.cmvr.inspection.domain.InspectionTaskInstance;
import com.cmvr.inspection.domain.vo.InspectionTaskStartVO;
import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo; import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo;
import com.cmvr.inspection.service.IInspectionTaskInstanceService; import com.cmvr.inspection.service.IInspectionTaskInstanceService;
import com.cmvr.common.utils.poi.ExcelUtil; import com.cmvr.common.utils.poi.ExcelUtil;
@ -121,9 +122,10 @@ public class InspectionTaskInstanceController extends BaseController
@PreAuthorize("@ss.hasPermi('inspection:taskInstance:edit')") @PreAuthorize("@ss.hasPermi('inspection:taskInstance:edit')")
@Log(title = "巡检任务执行实例", businessType = BusinessType.UPDATE) @Log(title = "巡检任务执行实例", businessType = BusinessType.UPDATE)
@PutMapping("/start/{id}") @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));
} }
/** /**

View File

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

View File

@ -6,6 +6,7 @@ import com.baomidou.mybatisplus.spring.service.IService;
import com.cmvr.aima.domain.AimaTaskInstance; import com.cmvr.aima.domain.AimaTaskInstance;
import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery;
import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo;
import com.cmvr.aima.domain.vo.AimaTaskStartVO;
/** /**
* 爱玛任务执行实例Service接口 * 爱玛任务执行实例Service接口
@ -28,7 +29,7 @@ public interface IAimaTaskInstanceService extends IService<AimaTaskInstance> {
* @param id 任务实例ID * @param id 任务实例ID
* @return 结果 * @return 结果
*/ */
int startInstance(String id); int startInstance(String id, AimaTaskStartVO startVO);
/** /**
* 暂停任务实例 * 暂停任务实例

View File

@ -6,6 +6,7 @@ import com.alibaba.fastjson2.JSONObject;
import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl; import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl;
import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery; import com.cmvr.aima.domain.dto.AimaTaskInstanceQuery;
import com.cmvr.aima.domain.vo.AimaTaskInstanceVo; import com.cmvr.aima.domain.vo.AimaTaskInstanceVo;
import com.cmvr.aima.domain.vo.AimaTaskStartVO;
import com.cmvr.common.exception.ServiceException; import com.cmvr.common.exception.ServiceException;
import com.cmvr.common.utils.SecurityUtils; import com.cmvr.common.utils.SecurityUtils;
import com.cmvr.test.flow.control.FlowControlService; import com.cmvr.test.flow.control.FlowControlService;
@ -119,7 +120,7 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
* 开始执行任务实例 * 开始执行任务实例
*/ */
@Override @Override
public int startInstance(String id) { public int startInstance(String id, AimaTaskStartVO startVO) {
AimaTaskInstance instance = aimaTaskInstanceMapper.selectAimaTaskInstanceById(id); AimaTaskInstance instance = aimaTaskInstanceMapper.selectAimaTaskInstanceById(id);
if (instance == null) { if (instance == null) {
throw new ServiceException("任务实例不存在"); throw new ServiceException("任务实例不存在");
@ -148,6 +149,8 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
TeTaskExecuteNormalVO.builder() TeTaskExecuteNormalVO.builder()
.taskId(aimaTask.getTaskConfigId()) .taskId(aimaTask.getTaskConfigId())
.moduleCode(FlowExecutionModuleCodes.AIMA) .moduleCode(FlowExecutionModuleCodes.AIMA)
.robotId(startVO == null ? null : startVO.getRobotId())
.runtimeAssignments(startVO == null ? null : startVO.getRuntimeAssignments())
.runParams(runParams) .runParams(runParams)
.build() .build()
); );

View File

@ -90,6 +90,9 @@ public class EdgeAgvServiceImpl implements EdgeAgvService {
edgeCommonVO.getRobotId(), AgvServiceGrpc.AgvServiceBlockingStub.class); edgeCommonVO.getRobotId(), AgvServiceGrpc.AgvServiceBlockingStub.class);
AgvCommand.AgvNavigateToStationCommand.Request request = AgvCommand.AgvNavigateToStationCommand.Request.newBuilder() AgvCommand.AgvNavigateToStationCommand.Request request = AgvCommand.AgvNavigateToStationCommand.Request.newBuilder()
.setHeader(EdgeCommonUtil.buildRequest(edgeCommonVO.getDeviceId())) .setHeader(EdgeCommonUtil.buildRequest(edgeCommonVO.getDeviceId()))
.setOptions(AgvUtils.AgvMotionOptions.newBuilder()
.setAsynchronous(true)
.build())
.setStationId(stationId) .setStationId(stationId)
.build(); .build();
executeGrpcCall(() -> stub.navigateToStation(request)); executeGrpcCall(() -> stub.navigateToStation(request));

View File

@ -639,6 +639,37 @@ public final class AgvServiceGrpc {
return getStopMappingMethod; 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 * 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); 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() { @java.lang.Override public final io.grpc.ServerServiceDefinition bindService() {
return io.grpc.ServerServiceDefinition.builder(getServiceDescriptor()) return io.grpc.ServerServiceDefinition.builder(getServiceDescriptor())
.addMethod( .addMethod(
@ -1035,6 +1076,13 @@ public final class AgvServiceGrpc {
cmvr.api.Common.CommandHeader.Request, cmvr.api.Common.CommandHeader.Request,
cmvr.api.Common.CommandHeader.Feedback>( cmvr.api.Common.CommandHeader.Feedback>(
this, METHODID_STOP_MAPPING))) 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(); .build();
} }
} }
@ -1278,6 +1326,17 @@ public final class AgvServiceGrpc {
io.grpc.stub.ClientCalls.asyncUnaryCall( io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getStopMappingMethod(), getCallOptions()), request, responseObserver); 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( return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getStopMappingMethod(), getCallOptions(), request); 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( return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getStopMappingMethod(), getCallOptions()), request); 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; 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_START_MAPPING = 17;
private static final int METHODID_STREAM_MAP = 18; private static final int METHODID_STREAM_MAP = 18;
private static final int METHODID_STOP_MAPPING = 19; private static final int METHODID_STOP_MAPPING = 19;
private static final int METHODID_TRANSLATE = 20;
private static final class MethodHandlers<Req, Resp> implements private static final class MethodHandlers<Req, Resp> implements
io.grpc.stub.ServerCalls.UnaryMethod<Req, Resp>, io.grpc.stub.ServerCalls.UnaryMethod<Req, Resp>,
@ -1849,6 +1930,10 @@ public final class AgvServiceGrpc {
serviceImpl.stopMapping((cmvr.api.Common.CommandHeader.Request) request, serviceImpl.stopMapping((cmvr.api.Common.CommandHeader.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.Common.CommandHeader.Feedback>) responseObserver); (io.grpc.stub.StreamObserver<cmvr.api.Common.CommandHeader.Feedback>) responseObserver);
break; break;
case METHODID_TRANSLATE:
serviceImpl.translate((cmvr.api.AgvCommand.AgvTranslateCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.AgvCommand.AgvTranslateCommand.Feedback>) responseObserver);
break;
default: default:
throw new AssertionError(); throw new AssertionError();
} }
@ -1930,6 +2015,7 @@ public final class AgvServiceGrpc {
.addMethod(getStartMappingMethod()) .addMethod(getStartMappingMethod())
.addMethod(getStreamMapMethod()) .addMethod(getStreamMapMethod())
.addMethod(getStopMappingMethod()) .addMethod(getStopMappingMethod())
.addMethod(getTranslateMethod())
.build(); .build();
} }
} }

View File

@ -25,7 +25,7 @@ public final class AgvServiceOuterClass {
java.lang.String[] descriptorData = { java.lang.String[] descriptorData = {
"\n\032cmvr/api/agv_service.proto\022\010cmvr.api\032\025" + "\n\032cmvr/api/agv_service.proto\022\010cmvr.api\032\025" +
"cmvr/api/common.proto\032\032cmvr/api/agv_comm" + "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" + "ate\022(.cmvr.api.AgvRuntimeStateCommand.Re" +
"quest\032).cmvr.api.AgvRuntimeStateCommand." + "quest\032).cmvr.api.AgvRuntimeStateCommand." +
"Feedback\022r\n\023getNavigationStatus\022,.cmvr.a" + "Feedback\022r\n\023getNavigationStatus\022,.cmvr.a" +
@ -72,7 +72,9 @@ public final class AgvServiceOuterClass {
"\032&.cmvr.api.AgvMapStreamCommand.Feedback" + "\032&.cmvr.api.AgvMapStreamCommand.Feedback" +
"0\001\022P\n\013stopMapping\022\037.cmvr.api.CommandHead" + "0\001\022P\n\013stopMapping\022\037.cmvr.api.CommandHead" +
"er.Request\032 .cmvr.api.CommandHeader.Feed" + "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 descriptor = com.google.protobuf.Descriptors.FileDescriptor
.internalBuildGeneratedFileFrom(descriptorData, .internalBuildGeneratedFileFrom(descriptorData,

View File

@ -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 { message AgvSetVelocityCommand {
// 请求体。 // 请求体。

View File

@ -69,4 +69,8 @@ service AgvService {
// 停止当前建图/扫图会话。 // 停止当前建图/扫图会话。
rpc stopMapping(CommandHeader.Request) returns (CommandHeader.Feedback); rpc stopMapping(CommandHeader.Request) returns (CommandHeader.Feedback);
// 按指定速度平移固定距离。成功返回仅表示控制器已接受命令。
rpc translate(AgvTranslateCommand.Request) returns (AgvTranslateCommand.Feedback);
} }

View File

@ -14,6 +14,29 @@ message AgvPose2d {
double theta = 3; 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 车体坐标系下的平面速度。 // AGV 车体坐标系下的平面速度。
message AgvVelocity { message AgvVelocity {
// 车体 X 方向线速度,单位:米/秒。 // 车体 X 方向线速度,单位:米/秒。
@ -24,6 +47,9 @@ message AgvVelocity {
double wz = 3; double wz = 3;
} }
// AGV 电池状态。 // AGV 电池状态。
message AgvBatteryState { message AgvBatteryState {
// 电量比例,范围:[0, 1],例如 0.8 表示 80%。 // 电量比例,范围:[0, 1],例如 0.8 表示 80%。

View File

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

View File

@ -4,6 +4,7 @@ import java.util.List;
import com.baomidou.mybatisplus.spring.service.IService; import com.baomidou.mybatisplus.spring.service.IService;
import com.cmvr.inspection.domain.InspectionTaskInstance; 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.dto.InspectionTaskInstanceQuery;
import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo; import com.cmvr.inspection.domain.vo.InspectionTaskInstanceVo;
/** /**
@ -61,7 +62,7 @@ public interface IInspectionTaskInstanceService extends IService<InspectionTaskI
* @param id 任务实例主键 * @param id 任务实例主键
* @return 结果 * @return 结果
*/ */
int startInstance(String id); int startInstance(String id, InspectionTaskStartVO startVO);
/** /**
* 暂停任务实例 * 暂停任务实例

View File

@ -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 + "]的类型与流程设备用途不匹配");
}
}
}

View File

@ -25,6 +25,7 @@ import org.springframework.stereotype.Service;
import com.cmvr.common.utils.DateUtils; import com.cmvr.common.utils.DateUtils;
import com.cmvr.inspection.mapper.InspectionTaskInstanceMapper; import com.cmvr.inspection.mapper.InspectionTaskInstanceMapper;
import com.cmvr.inspection.domain.InspectionTaskInstance; import com.cmvr.inspection.domain.InspectionTaskInstance;
import com.cmvr.inspection.domain.vo.InspectionTaskStartVO;
import com.cmvr.inspection.service.IInspectionTaskInstanceService; import com.cmvr.inspection.service.IInspectionTaskInstanceService;
/** /**
@ -143,7 +144,7 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
* 开始执行任务实例 * 开始执行任务实例
*/ */
@Override @Override
public int startInstance(String id) public int startInstance(String id, InspectionTaskStartVO startVO)
{ {
InspectionTaskInstance instance = inspectionTaskInstanceMapper.selectInspectionTaskInstanceById(id); InspectionTaskInstance instance = inspectionTaskInstanceMapper.selectInspectionTaskInstanceById(id);
if (instance == null) { if (instance == null) {
@ -169,7 +170,8 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
// 调用真正的执行 // 调用真正的执行
String insId = flowTaskRuntimeService.executeTask(TeTaskExecuteNormalVO.builder() String insId = flowTaskRuntimeService.executeTask(TeTaskExecuteNormalVO.builder()
.taskId(inspectionTask.getTaskConfigId()) .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) .moduleCode(FlowExecutionModuleCodes.INSPECTION)
.runParams(runParams) .runParams(runParams)
.build()); .build());

View File

@ -7,11 +7,13 @@ import com.cmvr.test.enums.RunModeEnum;
import com.cmvr.test.enums.TaskStatusEnum; import com.cmvr.test.enums.TaskStatusEnum;
import com.cmvr.test.enums.TerminalStatusEnum; import com.cmvr.test.enums.TerminalStatusEnum;
import com.cmvr.test.flow.builder.TaskKeyBuilder; import com.cmvr.test.flow.builder.TaskKeyBuilder;
import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO;
import lombok.Data; import lombok.Data;
import java.util.Collections; import java.util.Collections;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;
import java.util.ArrayList;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -47,6 +49,15 @@ public class TaskContext {
/** 稳定的机器人锁键。 */ /** 稳定的机器人锁键。 */
private String lockKey; private String lockKey;
/** 本次任务实际使用的全部机器人ID。 */
private List<String> robotIds = new ArrayList<>();
/** 全部机器人锁键,用于多机器人任务的原子占用和释放。 */
private List<String> lockKeys = new ArrayList<>();
/** 任务启动时冻结的机器人角色分配。 */
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
/** /**
* 当前运行状态(运行中/已完成/失败/暂停/终止) * 当前运行状态(运行中/已完成/失败/暂停/终止)
*/ */

View File

@ -1,6 +1,7 @@
package com.cmvr.test.flow.context; package com.cmvr.test.flow.context;
import com.alibaba.fastjson2.JSONObject; import com.alibaba.fastjson2.JSONObject;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.test.enums.TaskStatusEnum; import com.cmvr.test.enums.TaskStatusEnum;
import com.cmvr.test.enums.TerminalStatusEnum; import com.cmvr.test.enums.TerminalStatusEnum;
import lombok.Getter; import lombok.Getter;
@ -8,6 +9,8 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.util.List; import java.util.List;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -32,30 +35,41 @@ public class TaskContextManager {
/** /**
* 注册任务上下文(任务启动时调用) * 注册任务上下文(任务启动时调用)
*/ */
public void register(TaskContext context) { public synchronized void register(TaskContext context) {
instContextMap.put(context.getInstId(), context); List<String> lockKeys = effectiveLockKeys(context);
if (context.getLockKey() != null) { for (String lockKey : lockKeys) {
robotInstMap.put(context.getLockKey(), context.getInstId()); if (robotInstMap.containsKey(lockKey)) {
throw new GlobalException("当前任务中的机器人正在执行其他任务,请稍后再试");
}
} }
instContextMap.put(context.getInstId(), context);
lockKeys.forEach(lockKey -> robotInstMap.put(lockKey, context.getInstId()));
context.setTerminalStatus(TerminalStatusEnum.RUNNING); 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); TaskContext ctx = instContextMap.remove(instId);
if (ctx != null) { if (ctx != null) {
ctx.setTerminalStatus(TerminalStatusEnum.READY); ctx.setTerminalStatus(TerminalStatusEnum.READY);
ctx.setTerminalLocked(false); ctx.setTerminalLocked(false);
if (ctx.getLockKey() != null) { effectiveLockKeys(ctx).forEach(lockKey -> robotInstMap.remove(lockKey, instId));
robotInstMap.remove(ctx.getLockKey()); log.info("注销任务上下文:instId={}, robotIds={}", instId, ctx.getRobotIds());
}
log.info("注销任务上下文:instId={}, robotId={}", instId, ctx.getRobotId());
} }
} }
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());
}
/** /**
* 获取任务上下文 * 获取任务上下文
*/ */

View File

@ -22,6 +22,7 @@ import org.springframework.stereotype.Component;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Collections;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -55,7 +56,7 @@ public class FlowControlService {
throw new GlobalException("任务已终止,无需重复操作"); throw new GlobalException("任务已终止,无需重复操作");
} }
// 异步终止 // 异步终止
executor.execute(() -> edgeSystemService.stopAll(ctx.getRobotId())); stopAllRobots(ctx);
// 统一记录日志 + 设置上下文状态 + 数据库状态 // 统一记录日志 + 设置上下文状态 + 数据库状态
taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED); taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED);
@ -80,7 +81,7 @@ public class FlowControlService {
throw new GlobalException("任务已处于暂停状态"); throw new GlobalException("任务已处于暂停状态");
} }
// 异步停止终端 // 异步停止终端
executor.execute(() -> edgeSystemService.stopAll(ctx.getRobotId())); stopAllRobots(ctx);
ctx.setPaused(true); ctx.setPaused(true);
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING); ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
ctx.setStatus(TaskStatusEnum.PAUSED); ctx.setStatus(TaskStatusEnum.PAUSED);
@ -196,4 +197,14 @@ public class FlowControlService {
// 发布任务恢复事件 // 发布任务恢复事件
taskInstHolder.syncStatusAsResumed(instId); 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)));
}
} }

View File

@ -40,6 +40,7 @@ public class FlowItemExecutor {
private final TaskInstHolder taskInstHolder; private final TaskInstHolder taskInstHolder;
private final FlowExecutionChainBuilder chainBuilder; private final FlowExecutionChainBuilder chainBuilder;
private final TaskThreadRegistry taskThreadRegistry; private final TaskThreadRegistry taskThreadRegistry;
private final FlowRuntimeAssignmentResolver runtimeAssignmentResolver;
/** /**
* 执行主流程(每个检测项调用一次) * 执行主流程(每个检测项调用一次)
@ -226,7 +227,7 @@ public class FlowItemExecutor {
} }
@NotNull @NotNull
private static TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node, private TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node,
TaskNodeExecuteMessage rootMessage, String nodeId, String nodeName, TaskNodeExecuteMessage rootMessage, String nodeId, String nodeName,
List<Integer> iterations) { List<Integer> iterations) {
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
@ -245,6 +246,7 @@ public class FlowItemExecutor {
message.setLoopNum(rootMessage.getLoopNum()); message.setLoopNum(rootMessage.getLoopNum());
message.setLoopArray(rootMessage.getLoopArray()); message.setLoopArray(rootMessage.getLoopArray());
runtimeAssignmentResolver.resolveNode(node, message);
return message; return message;
} }

View File

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

View File

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

View File

@ -48,6 +48,7 @@ public class FlowTaskEngine {
message.setInstId(instId); message.setInstId(instId);
message.setTaskId(taskId); message.setTaskId(taskId);
message.setRobotId(robotId); message.setRobotId(robotId);
message.setRuntimeAssignments(context.getRuntimeAssignments());
message.setItemId(itemId); message.setItemId(itemId);
if (runMode.equals(RunModeEnum.VI_PROJECT)) { if (runMode.equals(RunModeEnum.VI_PROJECT)) {
message.setLoopDetail(item.getSchemeInfo()); message.setLoopDetail(item.getSchemeInfo());

View File

@ -18,6 +18,8 @@ import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO;
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO; import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
import com.cmvr.test.model.vo.TeTaskExecuteProjectVO; import com.cmvr.test.model.vo.TeTaskExecuteProjectVO;
import com.cmvr.test.model.vo.TeTaskExecuteTrailVO; 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.ITeDetectionItemService;
import com.cmvr.test.service.ITeTaskInstService; import com.cmvr.test.service.ITeTaskInstService;
import com.cmvr.test.service.ITeTaskOrchestrationService; import com.cmvr.test.service.ITeTaskOrchestrationService;
@ -28,6 +30,11 @@ import org.springframework.beans.factory.ObjectProvider;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import java.util.List; 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 TaskInstHolder taskInstHolder;
private final FlowTaskAsyncDispatcher flowTaskAsyncDispatcher; private final FlowTaskAsyncDispatcher flowTaskAsyncDispatcher;
private final ObjectProvider<RobotExecutionTargetResolver> targetResolverProvider; private final ObjectProvider<RobotExecutionTargetResolver> targetResolverProvider;
private final FlowRuntimeDefinitionValidator runtimeDefinitionValidator;
@Transactional(rollbackFor = Exception.class) @Transactional(rollbackFor = Exception.class)
public String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO) { public String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO) {
RobotExecutionTarget target = resolveTarget(taskExecuteNormalVO.getRobotId()); ResolvedRuntimeResources resources = resolveResources(
taskExecuteNormalVO.getRobotId(), taskExecuteNormalVO.getRuntimeAssignments());
return executeTaskInternal( return executeTaskInternal(
target, resources,
taskExecuteNormalVO.getTaskId(), taskExecuteNormalVO.getTaskId(),
taskExecuteNormalVO.getModuleCode(), taskExecuteNormalVO.getModuleCode(),
RunModeEnum.NORMAL, RunModeEnum.NORMAL,
@ -62,18 +71,15 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
@Transactional(rollbackFor = Exception.class) @Transactional(rollbackFor = Exception.class)
public String executeTrialTask(TeTaskExecuteTrailVO taskExecuteTrailVO) { public String executeTrialTask(TeTaskExecuteTrailVO taskExecuteTrailVO) {
JSONObject runParams = taskExecuteTrailVO.getRunParams(); 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 robotId = target.getRobotId();
String lockKey = target.lockKey(); assertResourcesAvailable(resources);
if (lockKey == null) {
throw new GlobalException("机器人ID不能为空");
}
if (taskInstHolder.isRobotLocked(lockKey)) {
throw new GlobalException("当前机器人正在执行其他任务,请稍后再试");
}
// 构建试运行任务定义(临时任务) // 构建试运行任务定义(临时任务)
String itemId = taskExecuteTrailVO.getItemId(); String itemId = taskExecuteTrailVO.getItemId();
runtimeDefinitionValidator.validate(taskExecuteTrailVO.getFlowData(), resources.assignments());
TeDetectItemDeployFlowVO deployFlowVO = new TeDetectItemDeployFlowVO(); TeDetectItemDeployFlowVO deployFlowVO = new TeDetectItemDeployFlowVO();
deployFlowVO.setId(itemId); deployFlowVO.setId(itemId);
@ -96,7 +102,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
List<TeQueryTaskDetailDTO> details = CollUtil.newArrayList(teQueryTaskDetailDTO); 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) // 异步调度执行(与正式任务一致,只是模式为 TRIAL)
flowTaskAsyncDispatcher.submit(instId, taskId, robotId, RunModeEnum.TRIAL, details); flowTaskAsyncDispatcher.submit(instId, taskId, robotId, RunModeEnum.TRIAL, details);
@ -116,9 +122,10 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
@Transactional(rollbackFor = Exception.class) @Transactional(rollbackFor = Exception.class)
public String executeProjectTask(TeTaskExecuteProjectVO taskExecuteProjectVO) { public String executeProjectTask(TeTaskExecuteProjectVO taskExecuteProjectVO) {
RobotExecutionTarget target = resolveTarget(taskExecuteProjectVO.getRobotId()); ResolvedRuntimeResources resources = resolveResources(
taskExecuteProjectVO.getRobotId(), taskExecuteProjectVO.getRuntimeAssignments());
return executeTaskInternal( return executeTaskInternal(
target, resources,
taskExecuteProjectVO.getProjectId(), taskExecuteProjectVO.getProjectId(),
null, null,
RunModeEnum.VI_PROJECT, RunModeEnum.VI_PROJECT,
@ -130,19 +137,18 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
@Transactional(rollbackFor = Exception.class) @Transactional(rollbackFor = Exception.class)
public String executeTaskInternal(String robotId, String taskId, String moduleCode, public String executeTaskInternal(String robotId, String taskId, String moduleCode,
RunModeEnum runMode, JSONObject runParams, Object originalVO) { RunModeEnum runMode, JSONObject runParams, Object originalVO) {
return executeTaskInternal(new RobotExecutionTarget(robotId), taskId, return executeTaskInternal(resolveResources(robotId, null), taskId,
moduleCode, runMode, runParams, originalVO); 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) { RunModeEnum runMode, JSONObject runParams, Object originalVO) {
RobotExecutionTarget target = resources.primaryTarget();
String robotId = target.getRobotId(); String robotId = target.getRobotId();
String lockKey = target.lockKey(); assertResourcesAvailable(resources);
if (lockKey != null && taskInstHolder.isRobotLocked(lockKey)) {
throw new GlobalException("当前机器人正在执行其他任务,请稍后再试");
}
// 查询任务详情(检测项信息) // 查询任务详情(检测项信息)
List<TeQueryTaskDetailDTO> details = taskOrchestrationService.queryTaskDetail(taskId); 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(); int count = (int) details.stream().filter(dto -> ObjUtil.equals(dto.getIsDeploy(), "1")).count();
if (count > 0) { if (count > 0) {
throw new GlobalException("当前任务存在未发布的检测项!"); throw new GlobalException("当前任务存在未发布的检测项!");
@ -153,7 +159,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
try { try {
// 注册上下文 // 注册上下文
registerTaskContext(instId, taskId, moduleCode, target, null, runMode, runParams); registerTaskContext(instId, taskId, moduleCode, resources, null, runMode, runParams);
// 异步提交执行 // 异步提交执行
flowTaskAsyncDispatcher.submit(instId, taskId, robotId, runMode, details); 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) { String itemId, RunModeEnum mode, JSONObject runParams) {
RobotExecutionTarget target = resources.primaryTarget();
TaskContext ctx = new TaskContext(); TaskContext ctx = new TaskContext();
ctx.setInstId(instId); ctx.setInstId(instId);
ctx.setTaskId(taskId); ctx.setTaskId(taskId);
@ -197,6 +204,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
ctx.setRunMode(mode); ctx.setRunMode(mode);
ctx.setRobotId(target.getRobotId()); ctx.setRobotId(target.getRobotId());
ctx.setLockKey(target.lockKey()); 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.setItemId(itemId);
ctx.setStatus(TaskStatusEnum.RUNNING); ctx.setStatus(TaskStatusEnum.RUNNING);
ctx.setTerminalLocked(true); ctx.setTerminalLocked(true);
@ -216,4 +226,58 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
// 纯逻辑工作流不调用边缘设备,可不绑定机器人。 // 纯逻辑工作流不调用边缘设备,可不绑定机器人。
return new RobotExecutionTarget(null); 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);
}
}
} }

View File

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

View File

@ -39,6 +39,16 @@ public class FlowInputPrepareInterceptor extends AbstractFlowMsgPreInterceptor {
// 调用已有的 prepare 方法 // 调用已有的 prepare 方法
JSONObject inputParams = paramPreparer.prepare(graph, node, message); 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); message.setInputParams(inputParams);
log.debug("[InputPrepare] 节点 {} 准备输入参数: {}", message.getNodeId(), inputParams); log.debug("[InputPrepare] 节点 {} 准备输入参数: {}", message.getNodeId(), inputParams);

View File

@ -3,6 +3,7 @@ package com.cmvr.test.flow.runtime.message;
import com.alibaba.fastjson2.JSONObject; import com.alibaba.fastjson2.JSONObject;
import com.cmvr.test.enums.ActionEnum; import com.cmvr.test.enums.ActionEnum;
import com.cmvr.test.flow.builder.FlowGraph; import com.cmvr.test.flow.builder.FlowGraph;
import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO;
import lombok.Data; import lombok.Data;
import java.util.ArrayList; import java.util.ArrayList;
@ -82,6 +83,18 @@ public class TaskNodeExecuteMessage {
*/ */
private String robotId; private String robotId;
/** 流程机器人角色标识。 */
private String roleKey;
/** 本次任务内被选中的角色成员。 */
private String roleMemberKey;
/** 冻结的本次执行角色分配。 */
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
/** 节点参数名到真实设备ID的解析结果。 */
private JSONObject resolvedDeviceBindings = new JSONObject();
/** /**
* 节点入参 * 节点入参
*/ */

View File

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

View File

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

View File

@ -10,6 +10,9 @@ import lombok.NoArgsConstructor;
import jakarta.validation.constraints.NotEmpty; import jakarta.validation.constraints.NotEmpty;
import java.util.ArrayList;
import java.util.List;
@Data @Data
@ApiModel("正常任务执行VO") @ApiModel("正常任务执行VO")
@Builder @Builder
@ -26,6 +29,11 @@ public class TeTaskExecuteNormalVO {
@ApiModelProperty("机器人永久唯一ID") @ApiModelProperty("机器人永久唯一ID")
private String robotId; private String robotId;
@ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId")
@Builder.Default
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
@ApiModelProperty(value = "运行时参数") @ApiModelProperty(value = "运行时参数")
@Builder.Default
private JSONObject runParams = new JSONObject(); private JSONObject runParams = new JSONObject();
} }

View File

@ -7,6 +7,9 @@ import lombok.Data;
import jakarta.validation.constraints.NotEmpty; import jakarta.validation.constraints.NotEmpty;
import java.util.ArrayList;
import java.util.List;
@Data @Data
@ApiModel("项目执行VO") @ApiModel("项目执行VO")
public class TeTaskExecuteProjectVO { public class TeTaskExecuteProjectVO {
@ -18,6 +21,9 @@ public class TeTaskExecuteProjectVO {
@ApiModelProperty("机器人永久唯一ID") @ApiModelProperty("机器人永久唯一ID")
private String robotId; private String robotId;
@ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId")
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
@ApiModelProperty(value = "运行时参数") @ApiModelProperty(value = "运行时参数")
private JSONObject runParams = new JSONObject(); private JSONObject runParams = new JSONObject();
} }

View File

@ -7,6 +7,9 @@ import lombok.Data;
import jakarta.validation.constraints.NotEmpty; import jakarta.validation.constraints.NotEmpty;
import java.util.ArrayList;
import java.util.List;
@Data @Data
@ApiModel("试运行任务执行VO") @ApiModel("试运行任务执行VO")
public class TeTaskExecuteTrailVO { public class TeTaskExecuteTrailVO {
@ -20,9 +23,11 @@ public class TeTaskExecuteTrailVO {
private String flowData; private String flowData;
@ApiModelProperty("机器人永久唯一ID") @ApiModelProperty("机器人永久唯一ID")
@NotEmpty(message = "机器人ID不能为空")
private String robotId; private String robotId;
@ApiModelProperty("本次执行的机器人角色分配;为空时兼容使用robotId")
private List<ResourceRuntimeAssignmentVO> runtimeAssignments = new ArrayList<>();
@ApiModelProperty(value = "运行时参数") @ApiModelProperty(value = "运行时参数")
private JSONObject runParams = new JSONObject(); private JSONObject runParams = new JSONObject();
} }

View File

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

View File

@ -61,7 +61,7 @@ CREATE TABLE `inspection_robot_device` (
`remark` varchar(500) DEFAULT NULL, `remark` varchar(500) DEFAULT NULL,
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
KEY `idx_robot_device_robot_id` (`robot_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`) KEY `idx_robot_device_online` (`robot_id`, `online_status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='机器人自动发现设备'; ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='机器人自动发现设备';

View File

@ -35,7 +35,7 @@ create table inspection_robot_device (
remark varchar(500) null, remark varchar(500) null,
primary key (id), primary key (id),
index idx_robot_device_robot_id (robot_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) index idx_robot_device_online (robot_id, online_status)
) comment='Devices discovered from robot QUIC heartbeat'; ) comment='Devices discovered from robot QUIC heartbeat';

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