diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java index e4bc5f7..6425ae6 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java @@ -118,12 +118,11 @@ public class TaskInstHolder { String action, String params, String message, List iterations) { nodeInstService.logFailed(instId, taskId, itemId, nodeId, nodeType, operate, action, params, message,iterations); - syncStatus(instId, TaskStatusEnum.FAILED); - - // 先获取上下文(在注销前) + // 终止和并发失败都可能先注销上下文,失败状态仍需落库。 TaskContext ctx = taskContextManager.get(instId); String moduleCode = ctx != null ? ctx.getModuleCode() : null; - + syncStatus(instId, TaskStatusEnum.FAILED); + // 注销任务上下文 taskContextManager.unregister(instId); @@ -141,22 +140,27 @@ public class TaskInstHolder { public void syncStatus(String instId, TaskStatusEnum status) { TaskContext ctx = taskContextManager.get(instId); - switch (status) { - case PAUSED: - ctx.setPaused(true); - ctx.setStatus(TaskStatusEnum.PAUSED); - ctx.setTerminalStatus(TerminalStatusEnum.PAUSED); - break; - case STOPPED: - ctx.setStopped(true); - ctx.setPaused(false); - ctx.setTerminalLocked(false); - break; + if (ctx != null) { + switch (status) { + case PAUSED: + ctx.setPaused(true); + ctx.setStatus(TaskStatusEnum.PAUSED); + ctx.setTerminalStatus(TerminalStatusEnum.PAUSED); + break; + case STOPPED: + ctx.setStopped(true); + ctx.setPaused(false); + ctx.setTerminalLocked(false); + break; + } + ctx.setStatus(status); + ctx.updateTimestamp(); } - ctx.setStatus(status); - ctx.updateTimestamp(); taskInstService.updateStatus(instId, status); + if (ctx == null) { + return; + } // 发布任务状态变更事件(暂停/终止/恢复) switch (status) { case PAUSED: diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandler.java index d16509f..1db9263 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandler.java @@ -1,6 +1,7 @@ package com.cmvr.test.flow.runtime.dispatcher; import com.cmvr.test.enums.ActionEnum; +import com.cmvr.common.exception.GlobalException; import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; import com.cmvr.test.flow.runtime.operator.NodeOperateHandler; @@ -20,10 +21,16 @@ public class FlowFunctionNodeHandler implements FlowNodeTypeHandler { @Override public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) { ActionEnum action = message.getAction(); + if (action == null) { + throw new GlobalException("节点动作不能为空"); + } NodeOperateHandler handler = operateHandlerMap.get(action.getOperate()); + if (handler == null) { + throw new GlobalException("不支持的节点操作类型: " + action.getOperate()); + } - TaskNodeExecuteResult execute = handler.execute(message); - return TaskNodeExecuteResult.success(execute.getOutputParams()); + // 暂停或终止时操作处理器允许返回 null,由执行器统一结束本次调度。 + return handler.execute(message); } } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java index de6a52a..261ffa1 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java @@ -58,7 +58,12 @@ public class FlowItemExecutor { executeNode(graph, graph.getStartNode(), startNodeId, startInputDefs, rootMessage, onFinished, new ArrayList<>()); // 调度后续节点(并发) - flowTaskScheduler.start(graph, taskInstHolder.getContext(rootMessage.getInstId()), executor, onFinished); + TaskContext context = taskInstHolder.getContext(rootMessage.getInstId()); + if (context == null || context.isStopped()) { + notifyFinished(onFinished, rootMessage.getInstId(), startNodeId); + return; + } + flowTaskScheduler.start(graph, context, executor, onFinished); } /** @@ -86,8 +91,12 @@ public class FlowItemExecutor { rootMessage, onFinished, iterations); log.info("当前执行-子图"); // 启动子图调度 - flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()), - executor, onFinished, iterations); + TaskContext context = taskInstHolder.getContext(rootMessage.getInstId()); + if (context == null || context.isStopped()) { + notifyFinished(onFinished, rootMessage.getInstId(), subStartNodeId); + return; + } + flowTaskScheduler.subStart(subGraph, context, executor, onFinished, iterations); } /** @@ -120,7 +129,7 @@ public class FlowItemExecutor { TaskNodeExecuteResult result = chainBuilder.build(interceptors, dispatchNodeHandler).apply(message); if (ObjUtil.isEmpty(result) || !result.isSuccess()) { -// log.warn("节点执行失败,中断调度:nodeId={}, status={}", nodeId, result != null ? result.isSuccess() : "null"); + handleUnsuccessfulResult(message, result, onFinished); return; } @@ -147,6 +156,20 @@ public class FlowItemExecutor { n -> executeNode(graph, n, startNodeId, startInputDefs, rootMessage, onFinished, iterations), iterations ); + } catch (RuntimeException e) { + TaskContext ctx = taskInstHolder.getContext(instId); + try { + if (ctx != null && !ctx.isPaused() && !ctx.isStopped()) { + taskInstHolder.markFailed( + instId, rootMessage.getTaskId(), itemId, nodeId, node.getNodeType().name(), + node.getAction() == null ? null : node.getAction().getOperate(), + node.getAction() == null ? null : node.getAction().getAction(), + null, "节点执行异常: " + e.getMessage(), iterations); + } + } finally { + notifyFinished(onFinished, instId, nodeId); + } + throw e; } finally { // 如果不是暂停全中断 TaskContext ctx = taskInstHolder.getContext(instId); @@ -156,6 +179,52 @@ public class FlowItemExecutor { } } + private void handleUnsuccessfulResult(TaskNodeExecuteMessage message, TaskNodeExecuteResult result, + Runnable onFinished) { + TaskContext ctx = taskInstHolder.getContext(message.getInstId()); + if (result == null) { + // 暂停需要保留等待关系以便恢复;终止会注销上下文,必须释放等待线程。 + if (ctx == null || ctx.isStopped()) { + notifyFinished(onFinished, message.getInstId(), message.getNodeId()); + } else if (!ctx.isPaused()) { + try { + taskInstHolder.markFailed( + message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(), + message.getNodeType(), + message.getAction() == null ? null : message.getAction().getOperate(), + message.getAction() == null ? null : message.getAction().getAction(), + message.getInputParams() == null ? null : message.getInputParams().toJSONString(), + "节点执行结果为空", message.getIterations()); + } finally { + notifyFinished(onFinished, message.getInstId(), message.getNodeId()); + } + } + return; + } + + try { + if (ctx != null && !ctx.isPaused() && !ctx.isStopped()) { + taskInstHolder.markFailed( + message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(), + message.getNodeType(), + message.getAction() == null ? null : message.getAction().getOperate(), + message.getAction() == null ? null : message.getAction().getAction(), + message.getInputParams() == null ? null : message.getInputParams().toJSONString(), + result.getErrorMsg(), message.getIterations()); + } + } finally { + notifyFinished(onFinished, message.getInstId(), message.getNodeId()); + } + } + + private void notifyFinished(Runnable onFinished, String instId, String nodeId) { + try { + onFinished.run(); + } catch (Exception e) { + log.error("释放流程等待器失败: instId={}, nodeId={}", instId, nodeId, e); + } + } + @NotNull private static TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node, TaskNodeExecuteMessage rootMessage, String nodeId, String nodeName, diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java index 3fee470..5542f49 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java @@ -22,45 +22,60 @@ public class FlowTaskEngine { private final FlowItemExecutor flowTaskExecutor; + private final FlowTaskScheduler flowTaskScheduler; + public void run(String instId, String taskId, String robotId, RunModeEnum runMode, List items) { - TaskContext context = taskInstHolder.getContext(instId); - context.setRunMode(runMode); - JSONObject runParams = context.getRunParams(); - - int totalItemCount = items.size(); - - for (int i = 0; i < totalItemCount; i++) { - TeQueryTaskDetailDTO item = items.get(i); - // 等待当前 item 完成后再执行下一个 - CountDownLatch latch = new CountDownLatch(1); - String itemId = item.getDetectItemId(); - context.setItemId(itemId); - - TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); - message.setInstId(instId); - message.setTaskId(taskId); - message.setRobotId(robotId); - message.setItemId(itemId); - if (runMode.equals(RunModeEnum.VI_PROJECT)) { - message.setLoopDetail(item.getSchemeInfo()); + try { + TaskContext context = taskInstHolder.getContext(instId); + if (context == null) { + log.info("任务已终止,跳过执行: instId={}", instId); + return; } + context.setRunMode(runMode); + JSONObject runParams = context.getRunParams(); - int pendingItemCount = totalItemCount - i - 1; - message.setPendingItemCount(pendingItemCount); + int totalItemCount = items.size(); - JSONObject jsonObject = runParams.getJSONObject(itemId); - if (!runMode.equals(RunModeEnum.TRIAL)) { - context.setRunParams(jsonObject); - } - flowTaskExecutor.executeGraph(message, item, latch::countDown); - try { - latch.await(); // 等待该检测项的子流程全部完成 - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - log.error("检测项执行被中断:itemId={}", itemId, e); - break; + for (int i = 0; i < totalItemCount; i++) { + TeQueryTaskDetailDTO item = items.get(i); + // 等待当前 item 完成后再执行下一个 + CountDownLatch latch = new CountDownLatch(1); + String itemId = item.getDetectItemId(); + context.setItemId(itemId); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setInstId(instId); + message.setTaskId(taskId); + message.setRobotId(robotId); + message.setItemId(itemId); + if (runMode.equals(RunModeEnum.VI_PROJECT)) { + message.setLoopDetail(item.getSchemeInfo()); + } + + int pendingItemCount = totalItemCount - i - 1; + message.setPendingItemCount(pendingItemCount); + + JSONObject jsonObject = runParams == null ? null : runParams.getJSONObject(itemId); + if (!runMode.equals(RunModeEnum.TRIAL)) { + context.setRunParams(jsonObject); + } + flowTaskExecutor.executeGraph(message, item, latch::countDown); + try { + latch.await(); // 等待该检测项的子流程全部完成 + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.info("检测项执行被终止:itemId={}", itemId); + break; + } + + context = taskInstHolder.getContext(instId); + if (context == null || context.isStopped()) { + break; + } } + } finally { + flowTaskScheduler.clear(instId); } } } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowAfterInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowAfterInterceptor.java index 304b142..5ecac18 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowAfterInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowAfterInterceptor.java @@ -28,6 +28,11 @@ public class FlowAfterInterceptor extends AbstractFlowMsgPreInterceptor{ @Override public TaskNodeExecuteResult doIntercept(TaskNodeExecuteMessage message, Function next) { TaskNodeExecuteResult result = next.apply(message); + TaskContext taskContext = taskInstHolder.getContext(message.getInstId()); + if (result == null && (taskContext == null || taskContext.isPaused() || taskContext.isStopped())) { + return null; + } + FlowExecutionEvent.EventType eventType; String errorMessage = null; @@ -38,7 +43,6 @@ public class FlowAfterInterceptor extends AbstractFlowMsgPreInterceptor{ errorMessage = result != null ? result.getErrorMsg() : "执行结果为空"; } - TaskContext taskContext = taskInstHolder.getContext(message.getInstId()); String moduleCode = taskContext != null ? taskContext.getModuleCode() : null; if (moduleCode == null) { diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowExceptionInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowExceptionInterceptor.java index 405a469..904e10b 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowExceptionInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowExceptionInterceptor.java @@ -4,6 +4,7 @@ import cn.hutool.core.exceptions.ExceptionUtil; import cn.hutool.core.util.ObjUtil; import cn.hutool.core.util.StrUtil; import com.cmvr.common.exception.GlobalException; +import com.cmvr.test.flow.context.TaskContext; import com.cmvr.test.flow.context.TaskInstHolder; import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher; import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; @@ -32,13 +33,20 @@ public class FlowExceptionInterceptor extends AbstractFlowMsgPreInterceptor { try { return next.apply(message); } catch (Exception e) { - log.error("节点执行异常:instId={}, nodeId={}, action={}", - message.getInstId(), message.getNodeId(), message.getAction(), e); - if (ExceptionUtil.getRootCause(e) instanceof InterruptedException) { - log.info("节点被中断退出: {}", message.getNodeId()); + log.info("节点因任务控制被中断: instId={}, nodeId={}", message.getInstId(), message.getNodeId()); return null; } + + TaskContext context = taskInstHolder.getContext(message.getInstId()); + if (context == null || context.isPaused() || context.isStopped()) { + log.info("任务已暂停或终止,忽略节点取消异常: instId={}, nodeId={}, error={}", + message.getInstId(), message.getNodeId(), ExceptionUtil.getMessage(e)); + return null; + } + + log.error("节点执行异常:instId={}, nodeId={}, action={}", + message.getInstId(), message.getNodeId(), message.getAction(), e); taskInstHolder.markFailed( message.getInstId(), message.getTaskId(), diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java index 9113185..f927a12 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java @@ -48,8 +48,13 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor { return null; } + if (result == null || !result.isSuccess()) { + return result; + } + taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), - result.getOutputParams().toJSONString(), StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations()); + result.getOutputParams() == null ? null : result.getOutputParams().toJSONString(), + StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations()); long end = System.currentTimeMillis(); log.info("[ {} ]节点执行结束, 参数[ {} ], 循环次数[ {} ], 耗时[ {} ]", message.getAction(), message.getInputParams(), message.getLoopNum(), end - start); diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandlerTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandlerTest.java new file mode 100644 index 0000000..135a554 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowFunctionNodeHandlerTest.java @@ -0,0 +1,42 @@ +package com.cmvr.test.flow.runtime.dispatcher; + +import com.cmvr.test.enums.ActionEnum; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; +import com.cmvr.test.flow.runtime.operator.NodeOperateHandler; +import org.junit.Test; + +import java.util.Collections; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; + +public class FlowFunctionNodeHandlerTest { + + @Test + public void stoppedNodeMayReturnNullWithoutCausingAnotherException() { + NodeOperateHandler stoppedHandler = message -> null; + FlowFunctionNodeHandler handler = new FlowFunctionNodeHandler( + Collections.singletonMap("EDGE", stoppedHandler)); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setAction(ActionEnum.CAMERA_START); + + assertNull(handler.handle(message)); + } + + @Test + public void operationFailureIsNotConvertedToSuccess() { + NodeOperateHandler failedHandler = message -> TaskNodeExecuteResult.failure("device failed"); + FlowFunctionNodeHandler handler = new FlowFunctionNodeHandler( + Collections.singletonMap("EDGE", failedHandler)); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setAction(ActionEnum.CAMERA_START); + + TaskNodeExecuteResult result = handler.handle(message); + assertFalse(result.isSuccess()); + assertEquals("device failed", result.getErrorMsg()); + } +}