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 0729953..6332a42 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 @@ -3,6 +3,7 @@ package com.cmvr.test.flow.runtime.engine; import cn.hutool.core.bean.BeanUtil; import cn.hutool.core.util.ObjUtil; import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; import com.cmvr.test.enums.NodeTypeEnum; import com.cmvr.test.flow.builder.FlowEdge; import com.cmvr.test.flow.builder.FlowGraph; @@ -260,8 +261,7 @@ public class FlowItemExecutor { FlowEdge matchedEdge = graph.getEdgeByFromAnchorId(matchedAnchorId); if (matchedEdge == null) { - log.warn("未找到匹配的分支出边,anchorId={}, nodeId={}", matchedAnchorId, nodeId); - return; + throw new GlobalException("分支节点未匹配到可执行路径,节点: {},锚点: {}", node.getNodeName(), matchedAnchorId); } // 移除未命中的路径 diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskAsyncDispatcher.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskAsyncDispatcher.java index 6ddeae7..55c8bfd 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskAsyncDispatcher.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskAsyncDispatcher.java @@ -1,6 +1,10 @@ package com.cmvr.test.flow.runtime.engine; +import cn.hutool.core.exceptions.ExceptionUtil; +import com.cmvr.test.enums.ActionEnum; import com.cmvr.test.enums.RunModeEnum; +import com.cmvr.test.flow.context.TaskContext; +import com.cmvr.test.flow.context.TaskInstHolder; import com.cmvr.test.model.dto.TeQueryTaskDetailDTO; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -22,6 +26,7 @@ public class FlowTaskAsyncDispatcher { private final ThreadPoolTaskExecutor executor; private final FlowTaskEngine taskEngine; + private final TaskInstHolder taskInstHolder; /** * 异步提交流程任务 @@ -42,8 +47,18 @@ public class FlowTaskAsyncDispatcher { log.info("异步任务执行结束:threadName={}, instId={}, taskId={}", Thread.currentThread().getName(), instId, taskId); } catch (Exception e) { log.error("任务执行异常:threadName={}, instId={}, taskId={}, 错误={}", Thread.currentThread().getName(), instId, taskId, e.getMessage(), e); - Thread.currentThread().interrupt(); // 保持中断状态 + if (ExceptionUtil.getRootCause(e) instanceof InterruptedException) { + Thread.currentThread().interrupt(); + return; + } + TaskContext context = taskInstHolder.getContext(instId); + if (context != null && !context.isStopped()) { + taskInstHolder.markFailed( + instId, taskId, context.getItemId(), null, null, + ActionEnum.NONE.getOperate(), ActionEnum.NONE.getAction(), null, + "任务执行异常: " + e.getMessage(), null); + } } }); } -} \ No newline at end of file +}