feat(task): 添加任务实例重新执行功能
- 在AimaTaskInstanceController中添加rerun接口支持爱玛任务重新执行 - 在InspectionTaskInstanceController中添加rerun接口支持巡检任务重新执行 - 在TeTaskInstController中添加rerun接口支持测试任务重新执行 - 实现AimaTaskInstanceService的rerunInstance方法处理爱玛任务重跑逻辑 - 实现InspectionTaskInstanceService的rerunInstance方法处理巡检任务重跑逻辑 - 扩展FlowTaskRuntimeService添加loadExecutionRequest和rerunTask方法 - 修改createAndSaveTaskInstance方法保存执行请求快照到任务实例配置中 - 添加单元测试验证执行请求快照的恢复和向后兼容功能
This commit is contained in:
parent
c4a9f29780
commit
da679363a7
@ -88,6 +88,14 @@ public class AimaTaskInstanceController extends BaseController {
|
||||
return toAjax(aimaTaskInstanceService.startInstance(id, startVO));
|
||||
}
|
||||
|
||||
@ApiOperation("再次执行已完成或失败的爱玛任务")
|
||||
@PreAuthorize("@ss.hasPermi('aima:taskinstance:edit')")
|
||||
@Log(title = "任务执行实例", businessType = BusinessType.INSERT)
|
||||
@PostMapping("/rerun/{id}")
|
||||
public AjaxResult rerun(@PathVariable("id") String id) {
|
||||
return success(aimaTaskInstanceService.rerunInstance(id));
|
||||
}
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
*/
|
||||
|
||||
@ -128,6 +128,15 @@ public class InspectionTaskInstanceController extends BaseController
|
||||
return toAjax(inspectionTaskInstanceService.startInstance(id, startVO));
|
||||
}
|
||||
|
||||
@ApiOperation("再次执行已完成或失败的巡检任务")
|
||||
@PreAuthorize("@ss.hasPermi('inspection:taskInstance:edit')")
|
||||
@Log(title = "巡检任务执行实例", businessType = BusinessType.INSERT)
|
||||
@PostMapping("/rerun/{id}")
|
||||
public AjaxResult rerun(@PathVariable("id") String id)
|
||||
{
|
||||
return success(inspectionTaskInstanceService.rerunInstance(id));
|
||||
}
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
*/
|
||||
|
||||
@ -6,6 +6,7 @@ import com.cmvr.test.model.domain.TeNodeInst;
|
||||
import com.cmvr.test.model.vo.TeQueryTaskInstVO;
|
||||
import com.cmvr.test.model.vo.TeTaskInstVO;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService;
|
||||
import io.swagger.annotations.Api;
|
||||
import io.swagger.annotations.ApiOperation;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
@ -14,6 +15,9 @@ import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RequestParam;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.bind.annotation.PathVariable;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import com.cmvr.common.core.domain.AjaxResult;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
@ -24,6 +28,7 @@ import java.util.List;
|
||||
public class TeTaskInstController extends BaseController {
|
||||
|
||||
private final ITeTaskInstService teTaskInstService;
|
||||
private final FlowTaskRuntimeService flowTaskRuntimeService;
|
||||
|
||||
@ApiOperation("查询任务实例列表")
|
||||
@PreAuthorize("@ss.hasPermi('test:inst:list')")
|
||||
@ -42,4 +47,11 @@ public class TeTaskInstController extends BaseController {
|
||||
return getDataTable(list);
|
||||
}
|
||||
|
||||
}
|
||||
@ApiOperation("再次执行已完成或失败的任务实例")
|
||||
@PreAuthorize("@ss.hasPermi('test:config:execute')")
|
||||
@PostMapping("/rerun/{id}")
|
||||
public AjaxResult rerun(@PathVariable("id") String id) {
|
||||
return success(flowTaskRuntimeService.rerunTask(id));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@ -31,6 +31,9 @@ public interface IAimaTaskInstanceService extends IService<AimaTaskInstance> {
|
||||
*/
|
||||
int startInstance(String id, AimaTaskStartVO startVO);
|
||||
|
||||
/** 按历史执行参数创建并启动一个新的爱玛任务实例。 */
|
||||
String rerunInstance(String id);
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
*
|
||||
|
||||
@ -188,6 +188,36 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
|
||||
return aimaTaskInstanceMapper.updateAimaTaskInstance(update);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String rerunInstance(String id) {
|
||||
AimaTaskInstance source = aimaTaskInstanceMapper.selectAimaTaskInstanceById(id);
|
||||
if (source == null) {
|
||||
throw new ServiceException("任务实例不存在");
|
||||
}
|
||||
if (AimaTaskStatusEnum.SUCCESS.getCode() != source.getStatus()
|
||||
&& AimaTaskStatusEnum.FAILED.getCode() != source.getStatus()) {
|
||||
throw new ServiceException("只有已完成或失败的任务才能再次执行");
|
||||
}
|
||||
if (StrUtil.isBlank(source.getTaskInsId())) {
|
||||
throw new ServiceException("历史任务没有执行参数快照,无法再次执行");
|
||||
}
|
||||
TeTaskExecuteNormalVO snapshot = flowTaskRuntimeService.loadExecutionRequest(source.getTaskInsId());
|
||||
AimaTaskInstance next = new AimaTaskInstance();
|
||||
next.setTaskId(source.getTaskId());
|
||||
next.setPhoneId(source.getPhoneId());
|
||||
next.setVehicleId(source.getVehicleId());
|
||||
next.setExecutionNumber(source.getExecutionNumber() == null ? 1 : source.getExecutionNumber() + 1);
|
||||
next.setRemark(source.getRemark());
|
||||
insertAimaTaskInstance(next);
|
||||
|
||||
AimaTaskStartVO startVO = new AimaTaskStartVO();
|
||||
startVO.setRobotId(snapshot.getRobotId());
|
||||
startVO.setRuntimeAssignments(snapshot.getRuntimeAssignments());
|
||||
startVO.setRunParams(snapshot.getRunParams());
|
||||
startInstance(next.getId(), startVO);
|
||||
return next.getId();
|
||||
}
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
*/
|
||||
|
||||
@ -64,6 +64,9 @@ public interface IInspectionTaskInstanceService extends IService<InspectionTaskI
|
||||
*/
|
||||
int startInstance(String id, InspectionTaskStartVO startVO);
|
||||
|
||||
/** 按历史执行参数创建并启动一个新的巡检任务实例。 */
|
||||
String rerunInstance(String id);
|
||||
|
||||
/**
|
||||
* 暂停任务实例
|
||||
*
|
||||
|
||||
@ -28,6 +28,7 @@ import com.cmvr.test.service.ITeTaskOrchestrationService;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.Lazy;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
import com.cmvr.common.utils.DateUtils;
|
||||
import com.cmvr.inspection.mapper.InspectionTaskInstanceMapper;
|
||||
import com.cmvr.inspection.domain.InspectionTaskInstance;
|
||||
@ -228,6 +229,35 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
|
||||
return inspectionTaskInstanceMapper.updateInspectionTaskInstance(update);
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public String rerunInstance(String id) {
|
||||
InspectionTaskInstance source = inspectionTaskInstanceMapper.selectInspectionTaskInstanceById(id);
|
||||
if (source == null) {
|
||||
throw new ServiceException("任务实例不存在");
|
||||
}
|
||||
if (TaskStatusEnum.SUCCESS.getCode() != source.getStatus()
|
||||
&& TaskStatusEnum.FAILED.getCode() != source.getStatus()) {
|
||||
throw new ServiceException("只有已完成或失败的任务才能再次执行");
|
||||
}
|
||||
if (StrUtil.isBlank(source.getTaskInsId())) {
|
||||
throw new ServiceException("历史任务没有执行参数快照,无法再次执行");
|
||||
}
|
||||
TeTaskExecuteNormalVO snapshot = flowTaskRuntimeService.loadExecutionRequest(source.getTaskInsId());
|
||||
InspectionTaskInstance next = new InspectionTaskInstance();
|
||||
next.setTaskId(source.getTaskId());
|
||||
next.setRobotId(source.getRobotId());
|
||||
next.setRemark(source.getRemark());
|
||||
insertInspectionTaskInstance(next);
|
||||
|
||||
InspectionTaskStartVO startVO = new InspectionTaskStartVO();
|
||||
startVO.setRobotId(snapshot.getRobotId());
|
||||
startVO.setRuntimeAssignments(snapshot.getRuntimeAssignments());
|
||||
startVO.setRunParams(snapshot.getRunParams());
|
||||
startInstance(next.getId(), startVO);
|
||||
return next.getId();
|
||||
}
|
||||
|
||||
private void normalizeRuntimeRobotIdentities(List<ResourceRuntimeAssignmentVO> assignments) {
|
||||
if (assignments == null) {
|
||||
return;
|
||||
|
||||
@ -70,6 +70,50 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
public TeTaskExecuteNormalVO loadExecutionRequest(String instId) {
|
||||
TeTaskInst source = taskInstService.getById(instId);
|
||||
if (source == null) {
|
||||
throw new GlobalException("任务实例不存在");
|
||||
}
|
||||
if (!RunModeEnum.NORMAL.name().equals(source.getRunMode())) {
|
||||
throw new GlobalException("只有正式任务实例支持再次执行");
|
||||
}
|
||||
if (StrUtil.isNotBlank(source.getConfig())) {
|
||||
try {
|
||||
TeTaskExecuteNormalVO request = JSON.parseObject(source.getConfig(), TeTaskExecuteNormalVO.class);
|
||||
if (request != null && StrUtil.isNotBlank(request.getTaskId())) {
|
||||
return request;
|
||||
}
|
||||
} catch (RuntimeException ignored) {
|
||||
// 旧版本的 config 可能不是执行请求快照,继续使用实例基础字段兼容重跑。
|
||||
}
|
||||
}
|
||||
return TeTaskExecuteNormalVO.builder()
|
||||
.taskId(source.getTaskId())
|
||||
.robotId(source.getRobotId())
|
||||
.build();
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public String rerunTask(String instId) {
|
||||
TeTaskInst source = taskInstService.getById(instId);
|
||||
assertRerunnable(source);
|
||||
return executeTask(loadExecutionRequest(instId));
|
||||
}
|
||||
|
||||
private void assertRerunnable(TeTaskInst source) {
|
||||
if (source == null) {
|
||||
throw new GlobalException("任务实例不存在");
|
||||
}
|
||||
String success = String.valueOf(TaskStatusEnum.SUCCESS.getCode());
|
||||
String failed = String.valueOf(TaskStatusEnum.FAILED.getCode());
|
||||
if (!success.equals(source.getStatus()) && !failed.equals(source.getStatus())) {
|
||||
throw new GlobalException("只有已完成或失败的任务才能再次执行");
|
||||
}
|
||||
}
|
||||
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public String executeTrialTask(TeTaskExecuteTrailVO taskExecuteTrailVO) {
|
||||
JSONObject runParams = taskExecuteTrailVO.getRunParams();
|
||||
@ -91,7 +135,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
String taskId = taskOrchestrationService.buildTrialTask(itemId);
|
||||
|
||||
// 创建并保存任务实例
|
||||
TeTaskInst instance = createAndSaveTaskInstance(taskId, target, RunModeEnum.TRIAL);
|
||||
TeTaskInst instance = createAndSaveTaskInstance(taskId, target, RunModeEnum.TRIAL, taskExecuteTrailVO);
|
||||
String instId = instance.getId();
|
||||
|
||||
try {
|
||||
@ -164,7 +208,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
throw new GlobalException("当前任务存在未发布的检测项!");
|
||||
}
|
||||
// 创建任务实例
|
||||
TeTaskInst instance = createAndSaveTaskInstance(taskId, target, runMode);
|
||||
TeTaskInst instance = createAndSaveTaskInstance(taskId, target, runMode, originalVO);
|
||||
String instId = instance.getId();
|
||||
|
||||
try {
|
||||
@ -188,10 +232,12 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
/**
|
||||
* 创建并保存任务实例
|
||||
*/
|
||||
private TeTaskInst createAndSaveTaskInstance(String taskId, RobotExecutionTarget target, RunModeEnum runMode) {
|
||||
private TeTaskInst createAndSaveTaskInstance(String taskId, RobotExecutionTarget target, RunModeEnum runMode,
|
||||
Object executionRequest) {
|
||||
TeTaskInst instance = new TeTaskInst();
|
||||
instance.setTaskId(taskId);
|
||||
instance.setRobotId(target.getRobotId());
|
||||
instance.setConfig(executionRequest == null ? null : JSON.toJSONString(executionRequest));
|
||||
instance.setRunMode(runMode.name());
|
||||
instance.setStatus(String.valueOf(TaskStatusEnum.RUNNING.getCode()));
|
||||
boolean save = taskInstService.save(instance);
|
||||
|
||||
@ -18,6 +18,15 @@ public interface FlowTaskRuntimeService {
|
||||
*/
|
||||
String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO);
|
||||
|
||||
/**
|
||||
* 读取正式任务实例保存的执行请求快照。
|
||||
* 旧实例没有快照时,使用实例中的任务和机器人信息生成兼容请求。
|
||||
*/
|
||||
TeTaskExecuteNormalVO loadExecutionRequest(String instId);
|
||||
|
||||
/** 使用历史实例的执行请求创建并启动一个新的正式任务实例。 */
|
||||
String rerunTask(String instId);
|
||||
|
||||
|
||||
/**
|
||||
* 执行试运行任务(TRIAL 模式)
|
||||
|
||||
@ -0,0 +1,71 @@
|
||||
package com.cmvr.test.flow.runtime.engine;
|
||||
|
||||
import com.alibaba.fastjson2.JSON;
|
||||
import com.alibaba.fastjson2.JSONObject;
|
||||
import com.cmvr.test.enums.RunModeEnum;
|
||||
import com.cmvr.test.model.domain.TeTaskInst;
|
||||
import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO;
|
||||
import com.cmvr.test.model.vo.RobotRuntimeMemberVO;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.List;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
public class FlowTaskRuntimeEntrySnapshotTest {
|
||||
|
||||
@Test
|
||||
public void restoresCompleteExecutionRequestSnapshot() {
|
||||
RobotRuntimeMemberVO member = new RobotRuntimeMemberVO();
|
||||
member.setRobotId("robot-2");
|
||||
ResourceRuntimeAssignmentVO assignment = new ResourceRuntimeAssignmentVO();
|
||||
assignment.setRoleKey("inspection");
|
||||
assignment.setMembers(List.of(member));
|
||||
TeTaskExecuteNormalVO request = TeTaskExecuteNormalVO.builder()
|
||||
.taskId("task-1")
|
||||
.robotId("robot-2")
|
||||
.runtimeAssignments(List.of(assignment))
|
||||
.runParams(new JSONObject().fluentPut("item-1:1",
|
||||
new JSONObject().fluentPut("threshold", 12)))
|
||||
.build();
|
||||
TeTaskInst source = source(JSON.toJSONString(request));
|
||||
|
||||
TeTaskExecuteNormalVO restored = entry(source).loadExecutionRequest(source.getId());
|
||||
|
||||
assertEquals("task-1", restored.getTaskId());
|
||||
assertEquals("robot-2", restored.getRobotId());
|
||||
assertEquals("inspection", restored.getRuntimeAssignments().get(0).getRoleKey());
|
||||
assertEquals(12, restored.getRunParams().getJSONObject("item-1:1").getIntValue("threshold"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void fallsBackForLegacyConfig() {
|
||||
TeTaskInst source = source("{\"legacyParam\":true}");
|
||||
|
||||
TeTaskExecuteNormalVO restored = entry(source).loadExecutionRequest(source.getId());
|
||||
|
||||
assertEquals(source.getTaskId(), restored.getTaskId());
|
||||
assertEquals(source.getRobotId(), restored.getRobotId());
|
||||
}
|
||||
|
||||
private FlowTaskRuntimeEntry entry(TeTaskInst source) {
|
||||
ITeTaskInstService service = (ITeTaskInstService) Proxy.newProxyInstance(
|
||||
ITeTaskInstService.class.getClassLoader(),
|
||||
new Class<?>[]{ITeTaskInstService.class},
|
||||
(proxy, method, args) -> "getById".equals(method.getName()) ? source : null);
|
||||
return new FlowTaskRuntimeEntry(null, service, null, null, null, null, null);
|
||||
}
|
||||
|
||||
private TeTaskInst source(String config) {
|
||||
TeTaskInst source = new TeTaskInst();
|
||||
source.setId("inst-1");
|
||||
source.setTaskId("legacy-task");
|
||||
source.setRobotId("legacy-robot");
|
||||
source.setRunMode(RunModeEnum.NORMAL.name());
|
||||
source.setConfig(config);
|
||||
return source;
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user