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 879252d..1edaae8 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 @@ -49,30 +49,68 @@ public class TaskInstHolder { nodeInstService.logRunning(instId, taskId, itemId, nodeId, nodeType, operate, action, paramsIn, iterations); } + /** + * 标记节点执行成功(仅记录节点日志,不更新任务状态) + * + * @param instId 任务实例ID + * @param itemId 检测项ID + * @param nodeId 节点ID + * @param pendingItemCount 剩余待执行检测项数量 + * @param nodeType 节点类型 + * @param paramsOut 输出参数 + * @param message 日志消息 + * @param iterations 迭代路径 + */ + public void markNodeSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType, + String paramsOut, String message, List iterations) { + // 仅记录节点执行日志 + nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message, iterations); + + // 如果是最后一个检测项的END节点,标记任务完成 + if (pendingItemCount == 0 && NodeTypeEnum.END.getCode().equalsIgnoreCase(nodeType)) { + markTaskCompleted(instId, itemId); + } + } + + /** + * 标记任务执行完成(更新任务状态为SUCCESS,发布完成事件) + * + * @param instId 任务实例ID + * @param itemId 检测项ID + */ + private void markTaskCompleted(String instId, String itemId) { + // 先获取上下文(在注销前) + TaskContext ctx = taskContextManager.get(instId); + String taskId = ctx != null ? ctx.getTaskId() : null; + String moduleCode = ctx != null ? ctx.getModuleCode() : null; + + // 更新任务状态为SUCCESS + syncStatus(instId, TaskStatusEnum.SUCCESS); + + // 注销任务上下文(释放终端锁) + taskContextManager.unregister(instId); + + // 发布任务完成事件 + eventPublisher.publishEvent(FlowExecutionEvent.builder() + .eventType(FlowExecutionEvent.EventType.TASK_COMPLETED) + .instId(instId) + .taskId(taskId) + .moduleCode(moduleCode) + .itemId(itemId) + .status(TaskStatusEnum.SUCCESS) + .build()); + + log.info("任务执行完成: instId={}, taskId={}", instId, taskId); + } + + /** + * 兼容旧方法名(已废弃,请使用 markNodeSuccess) + * @deprecated 使用 markNodeSuccess 替代 + */ + @Deprecated public void markSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType, String paramsOut, String message, List iterations) { - nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message,iterations); - syncStatus(instId, TaskStatusEnum.SUCCESS); - // 如果没有下一个检测项要执行 且是end节点 注销任务上下文 - if (pendingItemCount == 0 && nodeType.equalsIgnoreCase(NodeTypeEnum.END.getCode())) { - // 先获取上下文(在注销前) - TaskContext ctx = taskContextManager.get(instId); - String taskId = ctx != null ? ctx.getTaskId() : null; - String moduleCode = ctx != null ? ctx.getModuleCode() : null; - - // 注销任务上下文 - taskContextManager.unregister(instId); - - // 发布任务完成事件 - eventPublisher.publishEvent(FlowExecutionEvent.builder() - .eventType(FlowExecutionEvent.EventType.TASK_COMPLETED) - .instId(instId) - .taskId(taskId) - .moduleCode(moduleCode) - .itemId(itemId) - .status(TaskStatusEnum.SUCCESS) - .build()); - } + markNodeSuccess(instId, itemId, nodeId, pendingItemCount, nodeType, paramsOut, message, iterations); } public void markFailed(String instId, String taskId, String itemId, diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java index b14e31f..2669d4d 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java @@ -73,18 +73,21 @@ public class TaskThreadRegistry { */ public void interruptAll(String instId) { for (Map.Entry entry : nodeContextMap.entrySet()) { - if (entry.getKey().startsWith(instId)) { + if (entry.getKey().startsWith(instId + "_")) { FlowNodeWrapper node = entry.getValue().getNode(); if (node.getNodeType().equals(NodeTypeEnum.LOOP)) { - log.info("LOOP节点不直接中断"); + log.info("LOOP节点不直接中断,等待子图完成"); continue; } Thread thread = entry.getValue().getThread(); - log.info("中断线程:{}", thread.getName()); - thread.interrupt(); + if (thread != null && !thread.isInterrupted()) { + log.info("中断线程:{}", thread.getName()); + thread.interrupt(); + } } } -// releaseAllLatches(instId); + // ✅ 释放所有 latch,避免死锁 + releaseAllLatches(instId); log.info("-------------------------------------"); log.info("pending:\n{}", getPendingNodes(instId).stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList())); log.info("-------------------------------------"); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java index eabe6d6..f731ede 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java @@ -61,9 +61,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { JSONObject runParams = taskExecuteTrailVO.getRunParams(); String terminalId = taskExecuteTrailVO.getTerminalId(); if (StrUtil.isEmpty(terminalId)) { -// throw new GlobalException("终端ID不能为空"); + throw new GlobalException("终端ID不能为空"); } - if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) { + if (taskInstHolder.isTerminalLocked(terminalId)) { throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); } @@ -124,12 +124,13 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { @Transactional(rollbackFor = Exception.class) public String executeTaskInternal(String terminalId, String taskId, String moduleCode, RunModeEnum runMode, JSONObject runParams, Object originalVO) { + // 终端id可能为空 if (StrUtil.isEmpty(terminalId)) { // throw new GlobalException("终端ID不能为空"); } // 检查终端状态 - if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) { -// throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); + if (!StrUtil.isEmpty(terminalId) && taskInstHolder.isTerminalLocked(terminalId)) { + throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); } // 查询任务详情(检测项信息) List details = taskOrchestrationService.queryTaskDetail(taskId); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java index 91df724..5e47fd1 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java @@ -7,9 +7,9 @@ import lombok.Data; import lombok.NoArgsConstructor; /** - * 流程执行事件 + * 流程执行事件。 *

- * 携带流程执行过程中的关键信息,供外部模块监听和处理 + * 用于描述流程引擎执行过程中的关键节点和任务状态变化,供业务模块监听并处理。 */ @Data @Builder @@ -18,62 +18,62 @@ import lombok.NoArgsConstructor; public class FlowExecutionEvent { /** - * 事件类型 + * 事件类型。 */ private EventType eventType; /** - * 任务实例 ID + * 流程任务实例ID。 */ private String instId; /** - * 任务 ID + * 流程任务ID。 */ private String taskId; /** - * 模块编码 + * 事件所属业务模块编码,用于监听器路由隔离。 */ private String moduleCode; /** - * 节点数量 + * 当前检测项内的节点总数。 */ private Integer nodeCount; /** - * 当前检测项 ID + * 当前检测项ID。 */ private String itemId; /** - * 节点 ID(节点级事件时有值) + * 当前节点ID,节点级事件时有值。 */ private String nodeId; /** - * 节点类型(节点级事件时有值) + * 当前节点类型,节点级事件时有值。 */ private String nodeType; /** - * 节点名称 + * 当前节点名称。 */ private String nodeName; /** - * 任务状态 + * 当前任务状态。 */ private TaskStatusEnum status; /** - * 错误信息(失败事件时有值) + * 错误信息,失败事件时有值。 */ private String errorMessage; /** - * 事件类型枚举 + * 流程执行事件类型。 */ public enum EventType { /** 节点执行完成 */ @@ -84,9 +84,9 @@ public class FlowExecutionEvent { TASK_COMPLETED, /** 任务执行失败 */ TASK_FAILED, - /** 任务被暂停 */ + /** 任务暂停 */ TASK_PAUSED, - /** 任务被终止 */ + /** 任务终止 */ TASK_STOPPED, /** 任务恢复执行 */ TASK_RESUMED, 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 4fc24b7..9113185 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,7 +48,7 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor { return null; } - taskInstHolder.markSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), + taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), result.getOutputParams().toJSONString(), StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations()); long end = System.currentTimeMillis();