fix(flow): 修复任务终止逻辑并实现AGV导航等待功能
- 修复任务实例终止时对STOPPED状态的处理,添加PAUSED状态检查 - 实现AGV导航等待功能,支持超时控制和状态轮询 - 添加应用启动时遗留任务清理机制 - 优化任务终止时对空任务实例ID的检查 - 添加未完成节点记录的状态更新功能 - 实现AGV导航超时后的自动取消机制
This commit is contained in:
parent
eff8fdd4b6
commit
c4a1266335
@ -0,0 +1,77 @@
|
||||
package com.cmvr.web.config;
|
||||
|
||||
import com.cmvr.aima.domain.AimaTaskInstance;
|
||||
import com.cmvr.aima.enums.AimaTaskStatusEnum;
|
||||
import com.cmvr.aima.service.IAimaTaskInstanceService;
|
||||
import com.cmvr.inspection.domain.InspectionTaskInstance;
|
||||
import com.cmvr.inspection.service.IInspectionTaskInstanceService;
|
||||
import com.cmvr.test.enums.TaskStatusEnum;
|
||||
import com.cmvr.test.flow.control.FlowControlService;
|
||||
import com.cmvr.test.model.domain.TeTaskInst;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.boot.context.event.ApplicationReadyEvent;
|
||||
import org.springframework.context.event.EventListener;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
|
||||
@Slf4j
|
||||
@Component
|
||||
@RequiredArgsConstructor
|
||||
public class FlowTaskStartupRecovery {
|
||||
|
||||
private final ITeTaskInstService taskInstService;
|
||||
private final FlowControlService flowControlService;
|
||||
private final IAimaTaskInstanceService aimaTaskInstanceService;
|
||||
private final IInspectionTaskInstanceService inspectionTaskInstanceService;
|
||||
private final Date processStartedAt = new Date();
|
||||
|
||||
@EventListener(ApplicationReadyEvent.class)
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public void stopOrphanedTasks() {
|
||||
List<String> orphanedInstIds = taskInstService.lambdaQuery()
|
||||
.in(TeTaskInst::getStatus,
|
||||
String.valueOf(TaskStatusEnum.RUNNING.getCode()),
|
||||
String.valueOf(TaskStatusEnum.PAUSED.getCode()))
|
||||
.lt(TeTaskInst::getCreateTime, processStartedAt)
|
||||
.list()
|
||||
.stream()
|
||||
.map(TeTaskInst::getId)
|
||||
.toList();
|
||||
if (orphanedInstIds.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
|
||||
Date now = new Date();
|
||||
for (String instId : orphanedInstIds) {
|
||||
flowControlService.stopOrphaned(instId);
|
||||
}
|
||||
|
||||
aimaTaskInstanceService.lambdaUpdate()
|
||||
.in(AimaTaskInstance::getTaskInsId, orphanedInstIds)
|
||||
.in(AimaTaskInstance::getStatus,
|
||||
AimaTaskStatusEnum.RUNNING.getCode(), AimaTaskStatusEnum.PAUSED.getCode())
|
||||
.set(AimaTaskInstance::getStatus, AimaTaskStatusEnum.STOPPED.getCode())
|
||||
.set(AimaTaskInstance::getEndTime, now)
|
||||
.set(AimaTaskInstance::getUpdateTime, now)
|
||||
.update();
|
||||
|
||||
inspectionTaskInstanceService.lambdaUpdate()
|
||||
.in(InspectionTaskInstance::getTaskInsId, orphanedInstIds)
|
||||
.in(InspectionTaskInstance::getStatus,
|
||||
com.cmvr.inspection.enums.TaskStatusEnum.RUNNING.getCode(),
|
||||
com.cmvr.inspection.enums.TaskStatusEnum.PAUSED.getCode())
|
||||
.set(InspectionTaskInstance::getStatus,
|
||||
com.cmvr.inspection.enums.TaskStatusEnum.STOPPED.getCode())
|
||||
.set(InspectionTaskInstance::getEndTime, now)
|
||||
.set(InspectionTaskInstance::getUpdateTime, now)
|
||||
.update();
|
||||
|
||||
log.warn("应用启动时自动终止遗留流程任务: count={}", orphanedInstIds.size());
|
||||
log.debug("应用启动时自动终止遗留流程任务明细: instIds={}", orphanedInstIds);
|
||||
}
|
||||
}
|
||||
@ -248,8 +248,12 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
|
||||
if (instance == null) {
|
||||
throw new ServiceException("任务实例不存在");
|
||||
}
|
||||
if (AimaTaskStatusEnum.STOPPED.getCode() == instance.getStatus()) {
|
||||
return 1;
|
||||
}
|
||||
if (AimaTaskStatusEnum.NOT_STARTED.getCode() != instance.getStatus()
|
||||
&& AimaTaskStatusEnum.RUNNING.getCode() != instance.getStatus()) {
|
||||
&& AimaTaskStatusEnum.RUNNING.getCode() != instance.getStatus()
|
||||
&& AimaTaskStatusEnum.PAUSED.getCode() != instance.getStatus()) {
|
||||
throw new ServiceException("已完成或已取消的任务实例不能终止");
|
||||
}
|
||||
AimaTaskInstance update = new AimaTaskInstance();
|
||||
@ -258,7 +262,9 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl<AimaTaskInstanceMap
|
||||
update.setEndTime(DateUtils.getNowDate());
|
||||
update.setUpdateTime(DateUtils.getNowDate());
|
||||
update.setUpdateBy(SecurityUtils.getUsername());
|
||||
if (instance.getTaskInsId() != null && !instance.getTaskInsId().isBlank()) {
|
||||
flowControlService.stop(instance.getTaskInsId());
|
||||
}
|
||||
return aimaTaskInstanceMapper.updateAimaTaskInstance(update);
|
||||
}
|
||||
|
||||
|
||||
@ -339,7 +339,12 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
|
||||
if (instance == null) {
|
||||
throw new ServiceException("任务实例不存在");
|
||||
}
|
||||
if (TaskStatusEnum.NOT_STARTED.getCode() != instance.getStatus() && TaskStatusEnum.RUNNING.getCode() != instance.getStatus()) {
|
||||
if (TaskStatusEnum.STOPPED.getCode() == instance.getStatus()) {
|
||||
return 1;
|
||||
}
|
||||
if (TaskStatusEnum.NOT_STARTED.getCode() != instance.getStatus()
|
||||
&& TaskStatusEnum.RUNNING.getCode() != instance.getStatus()
|
||||
&& TaskStatusEnum.PAUSED.getCode() != instance.getStatus()) {
|
||||
throw new ServiceException("已完成或已取消的任务实例不能终止");
|
||||
}
|
||||
InspectionTaskInstance update = new InspectionTaskInstance();
|
||||
@ -348,7 +353,9 @@ public class InspectionTaskInstanceServiceImpl extends ServiceImpl<InspectionTas
|
||||
update.setEndTime(DateUtils.getNowDate());
|
||||
update.setUpdateTime(DateUtils.getNowDate());
|
||||
update.setUpdateBy(SecurityUtils.getUsername());
|
||||
if (instance.getTaskInsId() != null && !instance.getTaskInsId().isBlank()) {
|
||||
flowControlService.stop(instance.getTaskInsId());
|
||||
}
|
||||
return inspectionTaskInstanceMapper.updateInspectionTaskInstance(update);
|
||||
}
|
||||
|
||||
|
||||
@ -88,9 +88,11 @@ public class TaskThreadRegistry {
|
||||
}
|
||||
// ✅ 释放所有 latch,避免死锁
|
||||
releaseAllLatches(instId);
|
||||
log.info("-------------------------------------");
|
||||
log.info("pending:\n{}", getPendingNodes(instId).stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList()));
|
||||
log.info("-------------------------------------");
|
||||
List<TaskNodeExecuteContext> pendingNodes = getPendingNodes(instId);
|
||||
if (!pendingNodes.isEmpty()) {
|
||||
log.debug("任务中断后待退出节点: instId={}, nodeIds={}", instId,
|
||||
pendingNodes.stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList()));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@ -1,6 +1,8 @@
|
||||
package com.cmvr.test.flow.control;
|
||||
|
||||
import cn.hutool.core.bean.BeanUtil;
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import com.alibaba.fastjson2.JSON;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.edge.client.service.EdgeSystemService;
|
||||
import com.cmvr.test.enums.NodeTypeEnum;
|
||||
@ -14,6 +16,11 @@ import com.cmvr.test.flow.runtime.engine.FlowItemExecutor;
|
||||
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteContext;
|
||||
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage;
|
||||
import com.cmvr.test.service.ITeNodeInstService;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
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 lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
@ -23,6 +30,8 @@ import jakarta.annotation.Resource;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.stream.Collectors;
|
||||
@ -45,6 +54,7 @@ public class FlowControlService {
|
||||
private final TaskThreadRegistry taskThreadRegistry;
|
||||
private final FlowItemExecutor flowItemExecutor;
|
||||
private final EdgeSystemService edgeSystemService;
|
||||
private final ITeTaskInstService taskInstService;
|
||||
|
||||
/**
|
||||
* 终止流程
|
||||
@ -53,10 +63,11 @@ public class FlowControlService {
|
||||
|
||||
TaskContext ctx = taskInstHolder.getContext(instId);
|
||||
if (ctx == null) {
|
||||
throw new GlobalException("任务不存在");
|
||||
stopPersistedTaskWithoutContext(instId, true);
|
||||
return;
|
||||
}
|
||||
if (ctx.isStopped()) {
|
||||
throw new GlobalException("任务已终止,无需重复操作");
|
||||
return;
|
||||
}
|
||||
// 异步终止
|
||||
stopAllRobots(ctx);
|
||||
@ -72,6 +83,77 @@ public class FlowControlService {
|
||||
}
|
||||
}
|
||||
|
||||
/** Startup recovery entry: persistently closes an orphan without per-task warning logs. */
|
||||
public void stopOrphaned(String instId) {
|
||||
TaskContext ctx = taskInstHolder.getContext(instId);
|
||||
if (ctx != null) {
|
||||
stop(instId);
|
||||
return;
|
||||
}
|
||||
stopPersistedTaskWithoutContext(instId, false);
|
||||
}
|
||||
|
||||
private void stopPersistedTaskWithoutContext(String instId, boolean logDetail) {
|
||||
TeTaskInst persisted = taskInstService.getById(instId);
|
||||
if (persisted == null) {
|
||||
throw new GlobalException("任务不存在");
|
||||
}
|
||||
TaskStatusEnum status;
|
||||
try {
|
||||
status = TaskStatusEnum.fromCode(Integer.parseInt(persisted.getStatus()));
|
||||
} catch (NumberFormatException exception) {
|
||||
throw new GlobalException("任务状态异常,无法终止: " + persisted.getStatus());
|
||||
}
|
||||
if (TaskStatusEnum.STOPPED == status) {
|
||||
return;
|
||||
}
|
||||
if (TaskStatusEnum.SUCCESS == status || TaskStatusEnum.FAILED == status) {
|
||||
throw new GlobalException("任务已结束,无法终止");
|
||||
}
|
||||
|
||||
taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED);
|
||||
nodeInstService.stopUnfinished(instId, "服务重启后执行上下文已丢失,任务已终止");
|
||||
stopPersistedRobots(persisted);
|
||||
if (logDetail) {
|
||||
log.warn("终止缺少运行上下文的持久化任务: instId={}, previousStatus={}", instId, status);
|
||||
}
|
||||
}
|
||||
|
||||
private void stopPersistedRobots(TeTaskInst persisted) {
|
||||
Set<String> robotIds = new LinkedHashSet<>();
|
||||
if (StrUtil.isNotBlank(persisted.getRobotId())) {
|
||||
robotIds.add(persisted.getRobotId());
|
||||
}
|
||||
if (StrUtil.isNotBlank(persisted.getConfig())) {
|
||||
try {
|
||||
TeTaskExecuteNormalVO request = JSON.parseObject(persisted.getConfig(), TeTaskExecuteNormalVO.class);
|
||||
if (request != null && request.getRuntimeAssignments() != null) {
|
||||
for (ResourceRuntimeAssignmentVO assignment : request.getRuntimeAssignments()) {
|
||||
if (assignment == null || assignment.getMembers() == null) {
|
||||
continue;
|
||||
}
|
||||
for (RobotRuntimeMemberVO member : assignment.getMembers()) {
|
||||
if (member != null && StrUtil.isNotBlank(member.getRobotId())) {
|
||||
robotIds.add(member.getRobotId());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("解析遗留任务机器人分配失败,使用实例主机器人兜底: instId={}, error={}",
|
||||
persisted.getId(), exception.getMessage());
|
||||
}
|
||||
}
|
||||
robotIds.forEach(robotId -> executor.execute(() -> {
|
||||
try {
|
||||
edgeSystemService.stopAll(robotId);
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("终止遗留任务设备失败: instId={}, robotId={}, error={}",
|
||||
persisted.getId(), robotId, exception.getMessage());
|
||||
}
|
||||
}));
|
||||
}
|
||||
|
||||
/**
|
||||
* 暂停任务
|
||||
*/
|
||||
|
||||
@ -0,0 +1,91 @@
|
||||
package com.cmvr.test.flow.runtime.operator.edge;
|
||||
|
||||
import cmvr.msgs.AgvUtils;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.edge.client.model.EdgeCommonVO;
|
||||
import com.cmvr.edge.client.service.EdgeAgvService;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
@Slf4j
|
||||
public final class AgvNavigationWaiter {
|
||||
|
||||
public static final long DEFAULT_WAIT_TIMEOUT_MS = TimeUnit.MINUTES.toMillis(15);
|
||||
public static final long DEFAULT_POLL_INTERVAL_MS = 1000L;
|
||||
|
||||
// Values mirror cmvr-es AgvTaskState: None, Waiting, Running, Paused,
|
||||
// Completed, Failed, Canceled.
|
||||
private static final int STATE_NONE = 0;
|
||||
private static final int STATE_COMPLETED = 4;
|
||||
private static final int STATE_FAILED = 5;
|
||||
private static final int STATE_CANCELED = 6;
|
||||
|
||||
private AgvNavigationWaiter() {
|
||||
}
|
||||
|
||||
public static AgvUtils.AgvNavigationStatus navigateToStationAndWait(
|
||||
EdgeAgvService edgeAgvService,
|
||||
EdgeCommonVO common,
|
||||
String stationId,
|
||||
long waitTimeoutMs,
|
||||
long pollIntervalMs) {
|
||||
if (waitTimeoutMs <= 0) {
|
||||
throw new GlobalException("导航等待超时时间必须大于0");
|
||||
}
|
||||
if (pollIntervalMs <= 0 || pollIntervalMs > waitTimeoutMs) {
|
||||
throw new GlobalException("导航状态轮询间隔必须大于0且不能超过等待超时时间");
|
||||
}
|
||||
|
||||
edgeAgvService.navigateToStation(common, stationId);
|
||||
long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(waitTimeoutMs);
|
||||
AgvUtils.AgvNavigationStatus lastStatus = null;
|
||||
|
||||
while (System.nanoTime() < deadlineNanos) {
|
||||
lastStatus = edgeAgvService.getNavigationStatus(common);
|
||||
int state = lastStatus.getState();
|
||||
if (state == STATE_COMPLETED) {
|
||||
log.info("AGV导航完成,设备ID: {},站点ID: {}", common.getDeviceId(), stationId);
|
||||
return lastStatus;
|
||||
}
|
||||
if (state == STATE_FAILED) {
|
||||
throw navigationFailed("失败", lastStatus);
|
||||
}
|
||||
if (state == STATE_CANCELED) {
|
||||
throw navigationFailed("已取消", lastStatus);
|
||||
}
|
||||
|
||||
long remainingMs = TimeUnit.NANOSECONDS.toMillis(deadlineNanos - System.nanoTime());
|
||||
if (remainingMs <= 0) {
|
||||
break;
|
||||
}
|
||||
try {
|
||||
Thread.sleep(Math.min(pollIntervalMs, remainingMs));
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new GlobalException("导航等待被中断", e);
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
edgeAgvService.cancelNavigation(common);
|
||||
} catch (RuntimeException cancelError) {
|
||||
log.warn("AGV导航等待超时后取消失败,设备ID: {},原因: {}", common.getDeviceId(), cancelError.getMessage());
|
||||
}
|
||||
String stateMessage = lastStatus == null || lastStatus.getState() == STATE_NONE
|
||||
? "未获取到有效导航状态"
|
||||
: statusDescription(lastStatus);
|
||||
throw new GlobalException("AGV导航超时({}毫秒):{}", waitTimeoutMs, stateMessage);
|
||||
}
|
||||
|
||||
private static GlobalException navigationFailed(String stateName, AgvUtils.AgvNavigationStatus status) {
|
||||
return new GlobalException("AGV导航{}:{}", stateName, statusDescription(status));
|
||||
}
|
||||
|
||||
private static String statusDescription(AgvUtils.AgvNavigationStatus status) {
|
||||
String message = status.getMessage();
|
||||
return message == null || message.trim().isEmpty()
|
||||
? "状态=" + status.getState() + ",进度=" + status.getProgress()
|
||||
: message;
|
||||
}
|
||||
}
|
||||
@ -87,9 +87,15 @@ public class EdgeAgvOperateService implements EdgeOperateService {
|
||||
EdgeCommonVO edgeCommonVO = new EdgeCommonVO();
|
||||
edgeCommonVO.setRobotId(robotId);
|
||||
edgeCommonVO.setDeviceId(deviceId);
|
||||
edgeAgvService.navigateToStation(edgeCommonVO, stationId.trim());
|
||||
log.info("AGV移动到指定站点任务下发成功,设备ID: {},站点ID: {}", deviceId, stationId);
|
||||
break;
|
||||
long waitTimeoutMs = inputParams.getLongValue("waitTimeoutMs", AgvNavigationWaiter.DEFAULT_WAIT_TIMEOUT_MS);
|
||||
long pollIntervalMs = inputParams.getLongValue("pollIntervalMs", AgvNavigationWaiter.DEFAULT_POLL_INTERVAL_MS);
|
||||
AgvUtils.AgvNavigationStatus status = AgvNavigationWaiter.navigateToStationAndWait(
|
||||
edgeAgvService, edgeCommonVO, stationId.trim(), waitTimeoutMs, pollIntervalMs);
|
||||
return TaskNodeExecuteResult.success(new JSONObject()
|
||||
.fluentPut("state", status.getState())
|
||||
.fluentPut("progress", status.getProgress())
|
||||
.fluentPut("message", status.getMessage())
|
||||
.fluentPut("stationId", stationId.trim()));
|
||||
}
|
||||
|
||||
default:
|
||||
|
||||
@ -17,6 +17,7 @@ import com.cmvr.edge.client.service.EdgeAgvService;
|
||||
import com.cmvr.device.service.InspectionAlertListenService;
|
||||
import com.cmvr.llm.service.LLMAiAgentPlatformService;
|
||||
import com.cmvr.test.enums.ActionEnum;
|
||||
import com.cmvr.test.flow.runtime.operator.edge.AgvNavigationWaiter;
|
||||
import com.cmvr.test.flow.runtime.dispatcher.FlowCodeNodeHandler;
|
||||
import com.cmvr.test.flow.runtime.dispatcher.FlowHttpNodeHandler;
|
||||
import com.cmvr.test.flow.runtime.dispatcher.FlowSleepNodeHandler;
|
||||
@ -262,9 +263,12 @@ public class FlowActionExecutorService {
|
||||
if (StrUtil.isBlank(stationId)) {
|
||||
throw new GlobalException("站点ID(stationId)不能为空");
|
||||
}
|
||||
edgeAgvService.navigateToStation(edgeCommonVO, stationId.trim());
|
||||
log.info("AGV移动到指定站点任务下发成功,设备ID: {},站点ID: {}", deviceId, stationId);
|
||||
return "AGV移动到指定站点任务下发成功";
|
||||
long waitTimeoutMs = payload.getLongValue("waitTimeoutMs", AgvNavigationWaiter.DEFAULT_WAIT_TIMEOUT_MS);
|
||||
long pollIntervalMs = payload.getLongValue("pollIntervalMs", AgvNavigationWaiter.DEFAULT_POLL_INTERVAL_MS);
|
||||
AgvNavigationWaiter.navigateToStationAndWait(
|
||||
edgeAgvService, edgeCommonVO, stationId.trim(), waitTimeoutMs, pollIntervalMs);
|
||||
log.info("AGV移动到指定站点完成,设备ID: {},站点ID: {}", deviceId, stationId);
|
||||
return "AGV移动到指定站点完成";
|
||||
}
|
||||
|
||||
default:
|
||||
|
||||
@ -59,5 +59,8 @@ public interface ITeNodeInstService extends IService<TeNodeInst> {
|
||||
* 更新节点
|
||||
*/
|
||||
void update(String instId, String nodeId, TaskStatusEnum nodeStatus);
|
||||
|
||||
/** Marks unfinished node records as stopped when their in-memory execution is gone. */
|
||||
void stopUnfinished(String instId, String message);
|
||||
}
|
||||
|
||||
|
||||
@ -141,6 +141,17 @@ public class ITeNodeInstServiceImpl extends ServiceImpl<TeNodeInstMapper, TeNode
|
||||
.set(TeNodeInst::getStatus, nodeStatus.name());
|
||||
this.update(wrapper);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stopUnfinished(String instId, String message) {
|
||||
LambdaUpdateWrapper<TeNodeInst> wrapper = Wrappers.lambdaUpdate();
|
||||
wrapper.eq(TeNodeInst::getInstId, instId)
|
||||
.in(TeNodeInst::getStatus, TaskStatusEnum.RUNNING.name(), TaskStatusEnum.PAUSED.name())
|
||||
.set(TeNodeInst::getStatus, TaskStatusEnum.STOPPED.name())
|
||||
.set(TeNodeInst::getMessage, message)
|
||||
.set(TeNodeInst::getEndTime, System.currentTimeMillis());
|
||||
this.update(wrapper);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@ -0,0 +1,74 @@
|
||||
package com.cmvr.test.flow.control;
|
||||
|
||||
import com.cmvr.test.enums.TaskStatusEnum;
|
||||
import com.cmvr.test.flow.context.TaskContext;
|
||||
import com.cmvr.test.flow.context.TaskInstHolder;
|
||||
import com.cmvr.test.flow.context.TaskThreadRegistry;
|
||||
import com.cmvr.test.model.domain.TeTaskInst;
|
||||
import com.cmvr.test.service.ITeNodeInstService;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
public class FlowControlServiceTest {
|
||||
|
||||
@Test
|
||||
public void stopsPersistedRunningTaskWhenRuntimeContextWasLost() {
|
||||
AtomicBoolean unfinishedNodesStopped = new AtomicBoolean();
|
||||
AtomicReference<TaskStatusEnum> persistedStatus = new AtomicReference<>();
|
||||
ITeTaskInstService taskInstService = taskInstService("1");
|
||||
ITeNodeInstService nodeInstService = nodeInstService(unfinishedNodesStopped);
|
||||
TaskInstHolder holder = new TaskInstHolder(null, null, null, null) {
|
||||
@Override
|
||||
public TaskContext getContext(String instId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void syncStatus(String instId, TaskStatusEnum status) {
|
||||
persistedStatus.set(status);
|
||||
}
|
||||
};
|
||||
FlowControlService service = new FlowControlService(
|
||||
null, holder, nodeInstService, new TaskThreadRegistry(), null, null, taskInstService);
|
||||
|
||||
service.stop("inst-1");
|
||||
|
||||
assertEquals(TaskStatusEnum.STOPPED, persistedStatus.get());
|
||||
assertTrue(unfinishedNodesStopped.get());
|
||||
}
|
||||
|
||||
private ITeTaskInstService taskInstService(String status) {
|
||||
return (ITeTaskInstService) Proxy.newProxyInstance(
|
||||
ITeTaskInstService.class.getClassLoader(),
|
||||
new Class<?>[]{ITeTaskInstService.class},
|
||||
(proxy, method, args) -> {
|
||||
if ("getById".equals(method.getName())) {
|
||||
TeTaskInst task = new TeTaskInst();
|
||||
task.setId(String.valueOf(args[0]));
|
||||
task.setStatus(status);
|
||||
return task;
|
||||
}
|
||||
throw new UnsupportedOperationException(method.getName());
|
||||
});
|
||||
}
|
||||
|
||||
private ITeNodeInstService nodeInstService(AtomicBoolean stopped) {
|
||||
return (ITeNodeInstService) Proxy.newProxyInstance(
|
||||
ITeNodeInstService.class.getClassLoader(),
|
||||
new Class<?>[]{ITeNodeInstService.class},
|
||||
(proxy, method, args) -> {
|
||||
if ("stopUnfinished".equals(method.getName())) {
|
||||
stopped.set(true);
|
||||
return null;
|
||||
}
|
||||
throw new UnsupportedOperationException(method.getName());
|
||||
});
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,95 @@
|
||||
package com.cmvr.test.flow.runtime.operator.edge;
|
||||
|
||||
import cmvr.msgs.AgvUtils;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.edge.client.model.EdgeCommonVO;
|
||||
import com.cmvr.edge.client.service.EdgeAgvService;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.ArrayDeque;
|
||||
import java.util.Queue;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
public class AgvNavigationWaiterTest {
|
||||
|
||||
@Test
|
||||
public void waitsUntilNavigationCompletes() {
|
||||
Queue<Integer> states = new ArrayDeque<>();
|
||||
states.add(2);
|
||||
states.add(4);
|
||||
AtomicInteger queryCount = new AtomicInteger();
|
||||
EdgeAgvService service = service(states, queryCount, new AtomicBoolean());
|
||||
|
||||
AgvUtils.AgvNavigationStatus result = AgvNavigationWaiter.navigateToStationAndWait(
|
||||
service, common(), "station-A", 1000, 1);
|
||||
|
||||
assertEquals(4, result.getState());
|
||||
assertEquals(2, queryCount.get());
|
||||
}
|
||||
|
||||
@Test(expected = GlobalException.class)
|
||||
public void failsWhenEdgeReportsNavigationFailure() {
|
||||
Queue<Integer> states = new ArrayDeque<>();
|
||||
states.add(5);
|
||||
AgvNavigationWaiter.navigateToStationAndWait(
|
||||
service(states, new AtomicInteger(), new AtomicBoolean()), common(), "station-A", 1000, 1);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void cancelsNavigationAfterTimeout() {
|
||||
Queue<Integer> states = new ArrayDeque<>();
|
||||
states.add(2);
|
||||
AtomicBoolean canceled = new AtomicBoolean();
|
||||
try {
|
||||
AgvNavigationWaiter.navigateToStationAndWait(
|
||||
service(states, new AtomicInteger(), canceled), common(), "station-A", 5, 1);
|
||||
} catch (GlobalException expected) {
|
||||
assertTrue(canceled.get());
|
||||
return;
|
||||
}
|
||||
throw new AssertionError("Expected navigation timeout");
|
||||
}
|
||||
|
||||
private EdgeCommonVO common() {
|
||||
EdgeCommonVO common = new EdgeCommonVO();
|
||||
common.setRobotId("robot-1");
|
||||
common.setDeviceId("agv-1");
|
||||
return common;
|
||||
}
|
||||
|
||||
private EdgeAgvService service(Queue<Integer> states,
|
||||
AtomicInteger queryCount,
|
||||
AtomicBoolean canceled) {
|
||||
AtomicInteger lastState = new AtomicInteger(0);
|
||||
return (EdgeAgvService) Proxy.newProxyInstance(
|
||||
EdgeAgvService.class.getClassLoader(),
|
||||
new Class<?>[]{EdgeAgvService.class},
|
||||
(proxy, method, args) -> {
|
||||
switch (method.getName()) {
|
||||
case "navigateToStation":
|
||||
return null;
|
||||
case "getNavigationStatus":
|
||||
queryCount.incrementAndGet();
|
||||
Integer next = states.poll();
|
||||
if (next != null) {
|
||||
lastState.set(next);
|
||||
}
|
||||
return AgvUtils.AgvNavigationStatus.newBuilder()
|
||||
.setState(lastState.get())
|
||||
.setProgress(lastState.get() == 4 ? 1.0 : 0.5)
|
||||
.setMessage(lastState.get() == 5 ? "navigation failed" : "")
|
||||
.build();
|
||||
case "cancelNavigation":
|
||||
canceled.set(true);
|
||||
return null;
|
||||
default:
|
||||
throw new UnsupportedOperationException(method.getName());
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user