diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/config/FlowTaskStartupRecovery.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/config/FlowTaskStartupRecovery.java new file mode 100644 index 0000000..53e013b --- /dev/null +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/config/FlowTaskStartupRecovery.java @@ -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 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); + } +} diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java index d060902..35ea153 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/service/impl/AimaTaskInstanceServiceImpl.java @@ -248,8 +248,12 @@ public class AimaTaskInstanceServiceImpl extends ServiceImpl x.getNode().getNodeId()).collect(Collectors.toList())); - log.info("-------------------------------------"); + List pendingNodes = getPendingNodes(instId); + if (!pendingNodes.isEmpty()) { + log.debug("任务中断后待退出节点: instId={}, nodeIds={}", instId, + pendingNodes.stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList())); + } } /** diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java index 0e1788e..6cd0074 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java @@ -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 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()); + } + })); + } + /** * 暂停任务 */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiter.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiter.java new file mode 100644 index 0000000..0cf8bff --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiter.java @@ -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; + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeAgvOperateService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeAgvOperateService.java index 0772950..949b431 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeAgvOperateService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/edge/EdgeAgvOperateService.java @@ -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: diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java index 488edc4..c8da129 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java @@ -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: diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/service/ITeNodeInstService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/service/ITeNodeInstService.java index db9240a..f585cc8 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/service/ITeNodeInstService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/service/ITeNodeInstService.java @@ -59,5 +59,8 @@ public interface ITeNodeInstService extends IService { * 更新节点 */ 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); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/service/impl/ITeNodeInstServiceImpl.java b/cmvr-iot-test/src/main/java/com/cmvr/test/service/impl/ITeNodeInstServiceImpl.java index 2e6a9b3..1dfcc9f 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/service/impl/ITeNodeInstServiceImpl.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/service/impl/ITeNodeInstServiceImpl.java @@ -141,6 +141,17 @@ public class ITeNodeInstServiceImpl extends ServiceImpl 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); + } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowControlServiceTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowControlServiceTest.java new file mode 100644 index 0000000..62ee934 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/control/FlowControlServiceTest.java @@ -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 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()); + }); + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiterTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiterTest.java new file mode 100644 index 0000000..0b90e14 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/edge/AgvNavigationWaiterTest.java @@ -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 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 states = new ArrayDeque<>(); + states.add(5); + AgvNavigationWaiter.navigateToStationAndWait( + service(states, new AtomicInteger(), new AtomicBoolean()), common(), "station-A", 1000, 1); + } + + @Test + public void cancelsNavigationAfterTimeout() { + Queue 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 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()); + } + }); + } +}