diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java index bd1bbbb..2ac53c8 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java @@ -10,6 +10,7 @@ import com.cmvr.llm.service.LLMAiAgentPlatformService; import com.cmvr.test.flow.context.TaskContextManager; import com.cmvr.test.flow.control.FlowControlService; import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService; +import com.cmvr.test.flow.runtime.subflow.FlowSubFlowDefinitionService; import com.cmvr.test.flow.runtime.operator.edge.ti.TiTouchOperateService; import com.cmvr.test.model.vo.FlowActionRequestVO; import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO; @@ -67,6 +68,7 @@ public class TeFlowController extends BaseController { private final ITeDetectionItemService teDetectionItemService; private final ITeDetectionItemVersionService detectionItemVersionService; + private final FlowSubFlowDefinitionService subFlowDefinitionService; private final ITeNodeInstService nodeInstService; private final FlowTaskRuntimeService flowTaskRuntimeService; private final TaskContextManager taskContextManager; @@ -108,6 +110,12 @@ public class TeFlowController extends BaseController { return success(detectionItemVersionService.listVersions(itemId)); } + @ApiOperation("查询当前项目可用的已发布子流程") + @GetMapping("/subflows") + public AjaxResult subFlows(@RequestParam String currentItemId) { + return success(subFlowDefinitionService.listOptions(currentItemId)); + } + @ApiOperation("查看工作流发布版本") @GetMapping("/version/{versionId}") public AjaxResult version(@PathVariable String versionId) { diff --git a/cmvr-iot-admin/src/main/resources/application.yml b/cmvr-iot-admin/src/main/resources/application.yml index 2d17ba5..a37790c 100644 --- a/cmvr-iot-admin/src/main/resources/application.yml +++ b/cmvr-iot-admin/src/main/resources/application.yml @@ -132,6 +132,10 @@ cmvr: enabled: true initial-delay-ms: 3000 refresh-ms: 1000 + heartbeat-check-ms: 5000 + heartbeat-timeout-ms: ${CMVR_ROBOT_HEARTBEAT_TIMEOUT_MS:${cmvr.quic.heartbeat-timeout-ms:15000}} + # Comma-separated IDs for pure logical workflows. These IDs bypass QUIC online checks. + virtual-robot-ids: CN-CMVR-MBLRV1-CHANGAN-20260814-001 quic: enabled: true bind-host: 0.0.0.0 diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java index e3f5f31..923a707 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/listener/AimaFlowExecutionListener.java @@ -194,7 +194,7 @@ public class AimaFlowExecutionListener implements FlowExecutionListener { + (StrUtil.isBlank(nodeError) ? "" : ":" + StrUtil.sub(nodeError, 0, 300))) .actualResult(AimaLogExecutionContext.attach(logOutputParams, itemId, event.getItemOrder(), event.getItemOccurrence(), - event.getEventType().name(), event.getNodeType())) + event.getEventType().name(), event.getNodeType(), event.getIterations())) .screenshotUrl(imageUrl) .videoUrl(videoUrl) .errorStack(nodeError) diff --git a/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java index 64c0305..48c0c23 100644 --- a/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java +++ b/cmvr-iot-aima/src/main/java/com/cmvr/aima/support/AimaLogExecutionContext.java @@ -3,6 +3,8 @@ package com.cmvr.aima.support; import com.alibaba.fastjson2.JSONObject; import com.cmvr.aima.domain.AimaTestLog; +import java.util.List; + /** * Persists flow item identity inside the existing actual_result column. */ @@ -21,6 +23,14 @@ public final class AimaLogExecutionContext { public static String attach(JSONObject outputParams, String itemId, Integer itemOrder, Integer itemOccurrence, String eventType, String nodeType) { + return attach(outputParams, itemId, itemOrder, itemOccurrence, + eventType, nodeType, List.of()); + } + + public static String attach(JSONObject outputParams, String itemId, + Integer itemOrder, Integer itemOccurrence, + String eventType, String nodeType, + List iterations) { JSONObject result = outputParams == null ? new JSONObject() : JSONObject.parseObject(outputParams.toJSONString()); @@ -30,6 +40,7 @@ public final class AimaLogExecutionContext { context.put("itemOccurrence", itemOccurrence); context.put("eventType", eventType); context.put("nodeType", nodeType); + context.put("iterations", iterations == null ? List.of() : iterations); result.put(CONTEXT_KEY, context); return result.toJSONString(); } diff --git a/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java index 6e2b9e9..bef2a91 100644 --- a/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java +++ b/cmvr-iot-aima/src/test/java/com/cmvr/aima/support/AimaLogExecutionContextTest.java @@ -4,6 +4,8 @@ import com.alibaba.fastjson2.JSONObject; import com.cmvr.aima.domain.AimaTestLog; import org.junit.Test; +import java.util.List; + import static org.junit.Assert.assertEquals; public class AimaLogExecutionContextTest { @@ -21,4 +23,14 @@ public class AimaLogExecutionContextTest { assertEquals(Integer.valueOf(4), log.getItemOrder()); assertEquals(Integer.valueOf(2), log.getItemOccurrence()); } + + @Test + public void persistsLoopIterationPath() { + String actualResult = AimaLogExecutionContext.attach( + new JSONObject(), "case-1", 4, 2, + "NODE_COMPLETED", "LLM", List.of(2, 3)); + + JSONObject context = JSONObject.parseObject(actualResult).getJSONObject("_executionContext"); + assertEquals(List.of(2, 3), context.getList("iterations", Integer.class)); + } } diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/config/RobotStateProperties.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/config/RobotStateProperties.java new file mode 100644 index 0000000..d34010e --- /dev/null +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/config/RobotStateProperties.java @@ -0,0 +1,57 @@ +package com.cmvr.inspection.config; + +import org.apache.commons.lang3.StringUtils; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.stereotype.Component; + +import java.util.ArrayList; +import java.util.Date; +import java.util.List; +import java.util.Locale; + +@Component +@ConfigurationProperties(prefix = "cmvr.inspection.robot-state") +public class RobotStateProperties { + + private long heartbeatTimeoutMs = 15_000L; + private List virtualRobotIds = new ArrayList<>(); + + public long getHeartbeatTimeoutMs() { + return heartbeatTimeoutMs; + } + + public void setHeartbeatTimeoutMs(long heartbeatTimeoutMs) { + this.heartbeatTimeoutMs = heartbeatTimeoutMs; + } + + public List getVirtualRobotIds() { + return virtualRobotIds; + } + + public void setVirtualRobotIds(List virtualRobotIds) { + this.virtualRobotIds = virtualRobotIds == null ? new ArrayList<>() : virtualRobotIds; + } + + public Date heartbeatCutoff(Date now) { + long timeout = Math.max(1_000L, heartbeatTimeoutMs); + return new Date(now.getTime() - timeout); + } + + public boolean isVirtualRobot(String robotId) { + String normalized = normalize(robotId); + return normalized != null && normalizedVirtualRobotIds().contains(normalized); + } + + public List normalizedVirtualRobotIds() { + return virtualRobotIds.stream() + .map(RobotStateProperties::normalize) + .filter(StringUtils::isNotBlank) + .distinct() + .toList(); + } + + private static String normalize(String robotId) { + String value = StringUtils.trimToNull(robotId); + return value == null ? null : value.toLowerCase(Locale.ROOT); + } +} diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java index a35b18e..57b0374 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/listener/InspectionFlowExecutionListener.java @@ -181,6 +181,14 @@ public class InspectionFlowExecutionListener implements FlowExecutionListener { if (StrUtil.isNotBlank(nodeError)) { logOutputParams.put("errorMessage", nodeError); } + JSONObject executionContext = new JSONObject(); + executionContext.put("itemId", itemId); + executionContext.put("itemOrder", event.getItemOrder()); + executionContext.put("itemOccurrence", event.getItemOccurrence()); + executionContext.put("eventType", event.getEventType().name()); + executionContext.put("nodeType", event.getNodeType()); + executionContext.put("iterations", event.getIterations() == null ? List.of() : event.getIterations()); + logOutputParams.put("_executionContext", executionContext); String imageUrl = FlowMediaParamResolver.lastUrl(outputParams, "imageUrl"); String videoUrl = FlowMediaParamResolver.lastUrl(outputParams, "videoUrl"); String audioUrl = FlowMediaParamResolver.lastUrl(outputParams, "audioUrl"); diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/mapper/InspectionRobotMapper.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/mapper/InspectionRobotMapper.java index a17545b..5930801 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/mapper/InspectionRobotMapper.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/mapper/InspectionRobotMapper.java @@ -35,7 +35,15 @@ public interface InspectionRobotMapper extends MPJBaseMapper /** * Select connected robots that currently report one online AGV device. */ - List selectRuntimeSyncCandidates(); + List selectRuntimeSyncCandidates(@Param("heartbeatCutoff") Date heartbeatCutoff, + @Param("excludedRobotIds") List excludedRobotIds); + + List selectStaleConnectedRobots(@Param("heartbeatCutoff") Date heartbeatCutoff, + @Param("excludedRobotIds") List excludedRobotIds); + + int markHeartbeatExpiredOffline(@Param("id") String id, + @Param("heartbeatCutoff") Date heartbeatCutoff, + @Param("offlineTime") Date offlineTime); int archiveInspectionRobotByIds(@Param("ids") String[] ids, @Param("updateBy") String updateBy); diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotDeviceAssignmentValidator.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotDeviceAssignmentValidator.java index 7cc3bfc..8ed9448 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotDeviceAssignmentValidator.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotDeviceAssignmentValidator.java @@ -2,6 +2,7 @@ package com.cmvr.inspection.service; import com.cmvr.common.exception.GlobalException; import com.cmvr.inspection.mapper.InspectionRobotDeviceMapper; +import com.cmvr.inspection.config.RobotStateProperties; import com.cmvr.test.flow.runtime.engine.RobotDeviceAssignmentValidator; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; @@ -17,9 +18,13 @@ public class InspectionRobotDeviceAssignmentValidator implements RobotDeviceAssi "cannot start streaming while microphone recording is active"; private final InspectionRobotDeviceMapper deviceMapper; + private final RobotStateProperties robotStateProperties; @Override public void validate(String robotId, String deviceId, Integer expectedDeviceKind) { + if (robotStateProperties.isVirtualRobot(robotId)) { + throw new GlobalException("虚拟机器人不能绑定边缘设备:" + robotId); + } var devices = deviceMapper.selectByIdentity(robotId, deviceId); if (devices.isEmpty()) { throw new GlobalException("设备[" + deviceId + "]不属于机器人[" + robotId + "]"); diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotEndpointResolver.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotEndpointResolver.java index 37270d9..56cb5a4 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotEndpointResolver.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotEndpointResolver.java @@ -4,18 +4,21 @@ import com.cmvr.common.exception.GlobalException; import com.cmvr.common.robot.RobotEndpoint; import com.cmvr.common.robot.RobotEndpointResolver; import com.cmvr.inspection.domain.InspectionRobot; +import com.cmvr.inspection.config.RobotStateProperties; import com.cmvr.inspection.mapper.InspectionRobotMapper; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Service; import java.util.List; +import java.util.Date; @Service @RequiredArgsConstructor public class InspectionRobotEndpointResolver implements RobotEndpointResolver { private final InspectionRobotMapper robotMapper; + private final RobotStateProperties robotStateProperties; @Override public RobotEndpoint resolve(String robotId) { @@ -23,6 +26,9 @@ public class InspectionRobotEndpointResolver implements RobotEndpointResolver { if (normalized == null) { throw new GlobalException("机器人ID不能为空"); } + if (robotStateProperties.isVirtualRobot(normalized)) { + throw new GlobalException("虚拟机器人不能调用边缘设备服务:" + normalized); + } List robots = robotMapper.selectByRobotId(normalized); if (robots.isEmpty()) { throw new GlobalException("机器人不存在:" + normalized); @@ -37,6 +43,10 @@ public class InspectionRobotEndpointResolver implements RobotEndpointResolver { if (!"1".equals(robot.getConnectStatus())) { throw new GlobalException("机器人当前离线:" + normalized); } + Date heartbeatCutoff = robotStateProperties.heartbeatCutoff(new Date()); + if (robot.getLastHeartbeatTime() == null || robot.getLastHeartbeatTime().before(heartbeatCutoff)) { + throw new GlobalException("机器人QUIC心跳已超时:" + normalized); + } if (StringUtils.isBlank(robot.getIpAddress()) || robot.getPort() == null) { throw new GlobalException("机器人尚未上报gRPC地址:" + normalized); } diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotExecutionTargetResolver.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotExecutionTargetResolver.java index 76b37e9..183fbfd 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotExecutionTargetResolver.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/InspectionRobotExecutionTargetResolver.java @@ -2,6 +2,7 @@ package com.cmvr.inspection.service; import com.cmvr.test.flow.runtime.engine.RobotExecutionTarget; import com.cmvr.test.flow.runtime.engine.RobotExecutionTargetResolver; +import com.cmvr.inspection.config.RobotStateProperties; import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Service; @@ -11,11 +12,14 @@ import org.springframework.stereotype.Service; public class InspectionRobotExecutionTargetResolver implements RobotExecutionTargetResolver { private final InspectionRobotEndpointResolver endpointResolver; + private final RobotStateProperties robotStateProperties; @Override public RobotExecutionTarget resolve(String robotId) { String normalized = StringUtils.trimToNull(robotId); - endpointResolver.resolve(normalized); + if (!robotStateProperties.isVirtualRobot(normalized)) { + endpointResolver.resolve(normalized); + } return new RobotExecutionTarget(normalized); } } diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/RobotQuicStateService.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/RobotQuicStateService.java index 1276543..7ee2ff9 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/RobotQuicStateService.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/service/RobotQuicStateService.java @@ -2,6 +2,7 @@ package com.cmvr.inspection.service; import com.cmvr.common.utils.DateUtils; import com.cmvr.edge.client.manage.GrpcServiceManager; +import com.cmvr.inspection.config.RobotStateProperties; import com.cmvr.inspection.domain.InspectionRobot; import com.cmvr.inspection.domain.InspectionRobotDevice; import com.cmvr.inspection.mapper.InspectionRobotDeviceMapper; @@ -33,6 +34,7 @@ public class RobotQuicStateService { private final InspectionRobotMapper robotMapper; private final InspectionRobotDeviceMapper deviceMapper; private final GrpcServiceManager grpcServiceManager; + private final RobotStateProperties robotStateProperties; @Value("${cmvr.quic.grpc-host-source:observed-source}") private String grpcHostSource; @@ -46,6 +48,10 @@ public class RobotQuicStateService { node.getNodeId(), event.getType()); return; } + if (robotStateProperties.isVirtualRobot(robotId)) { + log.debug("Ignoring QUIC state for virtual robot, robotId={}, eventType={}", robotId, event.getType()); + return; + } List matches = robotMapper.selectByRobotId(robotId); if (matches.isEmpty()) { @@ -89,7 +95,10 @@ public class RobotQuicStateService { robot.setQuicBootId(node.getBootId()); robot.setObservedIp(node.getObservedSourceIp()); robot.setSoftwareVersion(node.getSoftwareVersion()); - robot.setLastHeartbeatTime(eventTime); + if (event.getType() == NodeEventType.NODE_EVENT_TYPE_REGISTERED + || event.getType() == NodeEventType.NODE_EVENT_TYPE_HEARTBEAT) { + robot.setLastHeartbeatTime(eventTime); + } robot.setUpdateTime(DateUtils.getNowDate()); if (offline || !hasEnabledAgv(node)) { robot.setStatus("1"); @@ -117,6 +126,25 @@ public class RobotQuicStateService { } } + @Transactional(rollbackFor = Exception.class) + public int expireStaleRobots() { + Date now = DateUtils.getNowDate(); + Date cutoff = robotStateProperties.heartbeatCutoff(now); + List excludedRobotIds = robotStateProperties.normalizedVirtualRobotIds(); + int expired = 0; + for (InspectionRobot robot : robotMapper.selectStaleConnectedRobots(cutoff, excludedRobotIds)) { + if (robotMapper.markHeartbeatExpiredOffline(robot.getId(), cutoff, now) == 0) { + continue; + } + deviceMapper.markAllOffline(robot.getRobotId(), now); + grpcServiceManager.invalidateRobot(robot.getRobotId()); + expired++; + log.warn("Robot heartbeat expired and was marked offline, robotId={}, lastHeartbeatTime={}, timeoutMs={}", + robot.getRobotId(), robot.getLastHeartbeatTime(), robotStateProperties.getHeartbeatTimeoutMs()); + } + return expired; + } + private void synchronizeDevices(String robotId, NodeSnapshot node, Date eventTime) { List reportedIds = new ArrayList<>(); for (ManagedDeviceStatus status : node.getDeviceManager().getDevicesList()) { diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotHeartbeatExpiryTask.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotHeartbeatExpiryTask.java new file mode 100644 index 0000000..aff4b34 --- /dev/null +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotHeartbeatExpiryTask.java @@ -0,0 +1,21 @@ +package com.cmvr.inspection.task; + +import com.cmvr.inspection.service.RobotQuicStateService; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(prefix = "cmvr.quic", name = "enabled", havingValue = "true") +public class RobotHeartbeatExpiryTask { + + private final RobotQuicStateService robotQuicStateService; + + @Scheduled(initialDelayString = "${cmvr.inspection.robot-state.initial-delay-ms:3000}", + fixedDelayString = "${cmvr.inspection.robot-state.heartbeat-check-ms:5000}") + public void expireStaleRobots() { + robotQuicStateService.expireStaleRobots(); + } +} diff --git a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotRuntimeStateRefreshTask.java b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotRuntimeStateRefreshTask.java index c9d213f..fa6fb51 100644 --- a/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotRuntimeStateRefreshTask.java +++ b/cmvr-iot-inspection/src/main/java/com/cmvr/inspection/task/RobotRuntimeStateRefreshTask.java @@ -1,6 +1,7 @@ package com.cmvr.inspection.task; import com.cmvr.inspection.domain.InspectionRobot; +import com.cmvr.inspection.config.RobotStateProperties; import com.cmvr.inspection.mapper.InspectionRobotMapper; import com.cmvr.inspection.service.IInspectionRobotService; import org.slf4j.Logger; @@ -12,6 +13,7 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.Set; +import java.util.Date; import java.util.concurrent.ConcurrentHashMap; /** @@ -27,22 +29,27 @@ public class RobotRuntimeStateRefreshTask private final InspectionRobotMapper inspectionRobotMapper; private final IInspectionRobotService inspectionRobotService; private final TaskExecutor taskExecutor; + private final RobotStateProperties robotStateProperties; private final Set syncingRobotIds = ConcurrentHashMap.newKeySet(); public RobotRuntimeStateRefreshTask(InspectionRobotMapper inspectionRobotMapper, IInspectionRobotService inspectionRobotService, - @Qualifier("threadPoolTaskExecutor") TaskExecutor taskExecutor) + @Qualifier("threadPoolTaskExecutor") TaskExecutor taskExecutor, + RobotStateProperties robotStateProperties) { this.inspectionRobotMapper = inspectionRobotMapper; this.inspectionRobotService = inspectionRobotService; this.taskExecutor = taskExecutor; + this.robotStateProperties = robotStateProperties; } @Scheduled(initialDelayString = "${cmvr.inspection.robot-state.initial-delay-ms:3000}", fixedDelayString = "${cmvr.inspection.robot-state.refresh-ms:5000}") public void refreshConnectedRobots() { - for (InspectionRobot robot : inspectionRobotMapper.selectRuntimeSyncCandidates()) { + Date cutoff = robotStateProperties.heartbeatCutoff(new Date()); + for (InspectionRobot robot : inspectionRobotMapper.selectRuntimeSyncCandidates( + cutoff, robotStateProperties.normalizedVirtualRobotIds())) { if (robot.getId() == null || !syncingRobotIds.add(robot.getId())) { continue; } diff --git a/cmvr-iot-inspection/src/main/resources/mapper/inspection/InspectionRobotMapper.xml b/cmvr-iot-inspection/src/main/resources/mapper/inspection/InspectionRobotMapper.xml index d7d27fb..5ef7c75 100644 --- a/cmvr-iot-inspection/src/main/resources/mapper/inspection/InspectionRobotMapper.xml +++ b/cmvr-iot-inspection/src/main/resources/mapper/inspection/InspectionRobotMapper.xml @@ -192,6 +192,14 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" from inspection_robot r where r.connect_status = '1' and coalesce(r.archived_status, '0') = '0' + and r.last_heartbeat_time is not null + and r.last_heartbeat_time >= #{heartbeatCutoff} + + and lower(r.robot_id) not in + + #{robotId} + + and (select count(*) from inspection_robot_device d where d.robot_id = r.robot_id @@ -200,6 +208,29 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" and d.online_status = '1') = 1 + + + + update inspection_robot + set connect_status = '0', status = '1', + update_by = 'heartbeat-timeout', update_time = #{offlineTime} + where id = #{id} + and connect_status = '1' + and (last_heartbeat_time is null or last_heartbeat_time < #{heartbeatCutoff}) + + update inspection_robot set archived_status = '1', connect_status = '0', update_time = sysdate(), update_by = #{updateBy} diff --git a/cmvr-iot-quic/src/main/java/com/cmvr/quic/gateway/core/NodeEventBroker.java b/cmvr-iot-quic/src/main/java/com/cmvr/quic/gateway/core/NodeEventBroker.java index 8a27666..c2a0368 100644 --- a/cmvr-iot-quic/src/main/java/com/cmvr/quic/gateway/core/NodeEventBroker.java +++ b/cmvr-iot-quic/src/main/java/com/cmvr/quic/gateway/core/NodeEventBroker.java @@ -12,6 +12,7 @@ import org.springframework.stereotype.Component; import java.util.ArrayList; import java.util.Collection; +import java.util.Locale; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -65,10 +66,10 @@ public final class NodeEventBroker implements QuicEventBus, AutoCloseable { } } - /** 同一 nodeId 重连后,旧连接迟到的心跳或离线事件不能覆盖新连接状态。 */ + /** 同一机器人重连后,旧连接迟到的心跳或离线事件不能覆盖新连接状态。 */ private boolean acceptNodeUpdate(final NodeEventType type, final NodeSnapshot snapshot) { final AtomicBoolean accepted = new AtomicBoolean(false); - nodes.compute(snapshot.getNodeId(), (nodeId, current) -> { + nodes.compute(nodeIdentity(snapshot), (identity, current) -> { if (type == NodeEventType.NODE_EVENT_TYPE_REGISTERED) { accepted.set(true); return snapshot; @@ -82,6 +83,14 @@ public final class NodeEventBroker implements QuicEventBus, AutoCloseable { return accepted.get(); } + private static String nodeIdentity(NodeSnapshot snapshot) { + String robotId = snapshot.getRobotId().trim(); + if (!robotId.isEmpty()) { + return "robot:" + robotId.toLowerCase(Locale.ROOT); + } + return "node:" + snapshot.getNodeId(); + } + @Override public Subscription watchNodes(boolean includeCurrent, Consumer listener) { if (closed.get()) { diff --git a/cmvr-iot-quic/src/test/java/com/cmvr/quic/gateway/core/NodeEventBrokerTest.java b/cmvr-iot-quic/src/test/java/com/cmvr/quic/gateway/core/NodeEventBrokerTest.java index d70502a..9a6df83 100644 --- a/cmvr-iot-quic/src/test/java/com/cmvr/quic/gateway/core/NodeEventBrokerTest.java +++ b/cmvr-iot-quic/src/test/java/com/cmvr/quic/gateway/core/NodeEventBrokerTest.java @@ -34,6 +34,40 @@ public class NodeEventBrokerTest { } } + @Test + public void shouldKeepDifferentRobotsThatShareNodeId() { + NodeEventBroker broker = new NodeEventBroker(4); + try { + broker.publishNode(NodeEventType.NODE_EVENT_TYPE_REGISTERED, + node("edge-default", "robot-1", "session-1", true)); + broker.publishNode(NodeEventType.NODE_EVENT_TYPE_REGISTERED, + node("edge-default", "robot-2", "session-2", true)); + + Assert.assertEquals(2, broker.listNodes().size()); + } finally { + broker.close(); + } + } + + @Test + public void shouldRejectOldSessionWhenSameRobotChangesNodeId() { + NodeEventBroker broker = new NodeEventBroker(4); + try { + broker.publishNode(NodeEventType.NODE_EVENT_TYPE_REGISTERED, + node("edge-old", "robot-1", "old-session", true)); + broker.publishNode(NodeEventType.NODE_EVENT_TYPE_REGISTERED, + node("edge-new", "robot-1", "new-session", true)); + broker.publishNode(NodeEventType.NODE_EVENT_TYPE_OFFLINE, + node("edge-old", "robot-1", "old-session", false)); + + NodeSnapshot current = broker.listNodes().iterator().next(); + Assert.assertEquals("new-session", current.getSessionId()); + Assert.assertTrue(current.getOnline()); + } finally { + broker.close(); + } + } + @Test public void shouldPublishNodeEventsDirectlyAndAllowUnsubscribe() { NodeEventBroker broker = new NodeEventBroker(4); @@ -78,8 +112,13 @@ public class NodeEventBrokerTest { } private static NodeSnapshot node(String nodeId, String sessionId, boolean online) { + return node(nodeId, "", sessionId, online); + } + + private static NodeSnapshot node(String nodeId, String robotId, String sessionId, boolean online) { return NodeSnapshot.newBuilder() .setNodeId(nodeId) + .setRobotId(robotId) .setSessionId(sessionId) .setOnline(online) .build(); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java index 5703299..0a00edd 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java @@ -22,6 +22,7 @@ public enum ActionEnum { SUB_START("NONE", "SUB_START", "循环开始(子流程开始)"), END("NONE", "END", "结束"), START_LOOP("NONE", "START_LOOP", "开始循环"), + SUB_FLOW("NONE", "SUB_FLOW", "执行子流程"), STOP_LOOP("NONE", "STOP_LOOP", "结束循环"), BRANCH("NONE", "BRANCH", "分支"), SUB_END("NONE", "SUB_END", "子流程结束"), diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/NodeTypeEnum.java b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/NodeTypeEnum.java index 9296166..b503167 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/NodeTypeEnum.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/NodeTypeEnum.java @@ -15,6 +15,7 @@ public enum NodeTypeEnum { FUNCTION("function"), BRANCH("branch"), LOOP("loop"), + SUB_FLOW("subFlow"), STOP_LOOP("stopLoop"), END("end"), SLEEP("sleep"), diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowGraph.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowGraph.java index a1bd18d..5abfa08 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowGraph.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowGraph.java @@ -114,15 +114,25 @@ public class FlowGraph { * 获取子图中起始节点(入度为0) */ public FlowNodeWrapper getInDegreeZeroNode() { + List nodes = getInDegreeZeroNodes(); + return nodes.isEmpty() ? null : nodes.get(0); + } + + /** + * Returns every executable entry of a subgraph. A loop body can contain + * multiple independent lines, so selecting only the first entry loses work. + */ + public List getInDegreeZeroNodes() { lock.readLock().lock(); try { + List result = new ArrayList<>(); for (FlowNodeWrapper node : nodeMap.values()) { List flowEdges = predecessors.get(node.getNodeId()); if (flowEdges == null || flowEdges.isEmpty()) { - return node; + result.add(node); } } - return null; + return result; } finally { lock.readLock().unlock(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowModelBuilder.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowModelBuilder.java index f67f76e..6ec4bb8 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowModelBuilder.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/builder/FlowModelBuilder.java @@ -2,6 +2,7 @@ package com.cmvr.test.flow.builder; import cn.hutool.core.collection.CollUtil; import cn.hutool.core.util.ObjUtil; +import cn.hutool.core.util.StrUtil; import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSONArray; import com.alibaba.fastjson2.JSONObject; @@ -14,9 +15,11 @@ import com.cmvr.test.enums.ParamScope; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Set; /** * 构建任务执行图(支持主流程 + 子图 + 循环) @@ -30,6 +33,7 @@ public class FlowModelBuilder { JSONObject flow = JSON.parseObject(flowJson); JSONArray nodeArray = flow.getJSONArray("nodes"); JSONArray edgeArray = flow.getJSONArray("edges"); + validateNoCircularParameterReferences(nodeArray); FlowGraph graph = new FlowGraph(); @@ -57,6 +61,75 @@ public class FlowModelBuilder { return graph; } + private static void validateNoCircularParameterReferences(JSONArray nodes) { + if (nodes == null || nodes.isEmpty()) return; + Map> dependencies = new LinkedHashMap<>(); + Map nodeNames = new HashMap<>(); + for (int i = 0; i < nodes.size(); i++) { + JSONObject node = nodes.getJSONObject(i); + String nodeId = node.getString("id"); + JSONObject properties = node.getJSONObject("properties"); + Set references = new HashSet<>(); + collectActiveReferences(properties, references); + dependencies.put(nodeId, references); + nodeNames.put(nodeId, properties == null + ? nodeId : StrUtil.blankToDefault(properties.getString("name"), nodeId)); + } + + Map states = new HashMap<>(); + List path = new ArrayList<>(); + for (String nodeId : dependencies.keySet()) { + List cycle = findReferenceCycle(nodeId, dependencies, states, path); + if (cycle == null) continue; + List names = cycle.stream().map(id -> nodeNames.getOrDefault(id, id)).toList(); + throw new GlobalException("节点参数存在循环引用:" + String.join(" -> ", names)); + } + } + + private static void collectActiveReferences(Object value, Set references) { + if (value instanceof JSONArray array) { + for (Object child : array) collectActiveReferences(child, references); + return; + } + if (!(value instanceof JSONObject object)) return; + if ("quote".equals(object.getString("type"))) { + addReference(object.getJSONArray("quote"), references); + } + if ("quote".equals(object.getString("nameType"))) { + addReference(object.getJSONArray("nameQuote"), references); + } + for (Object child : object.values()) collectActiveReferences(child, references); + } + + private static void addReference(JSONArray quote, Set references) { + if (quote != null && quote.size() >= 3 && StrUtil.isNotBlank(quote.getString(0))) { + references.add(quote.getString(0)); + } + } + + private static List findReferenceCycle(String nodeId, + Map> dependencies, + Map states, + List path) { + if (Integer.valueOf(2).equals(states.get(nodeId))) return null; + if (Integer.valueOf(1).equals(states.get(nodeId))) { + int start = path.indexOf(nodeId); + List cycle = new ArrayList<>(path.subList(Math.max(0, start), path.size())); + cycle.add(nodeId); + return cycle; + } + states.put(nodeId, 1); + path.add(nodeId); + for (String sourceId : dependencies.getOrDefault(nodeId, Collections.emptySet())) { + if (!dependencies.containsKey(sourceId)) continue; + List cycle = findReferenceCycle(sourceId, dependencies, states, path); + if (cycle != null) return cycle; + } + path.remove(path.size() - 1); + states.put(nodeId, 2); + return null; + } + private static void setParams(Map nodeMap, FlowGraph graph) { Map> nodeParams = new LinkedHashMap<>(); for (Map.Entry entry : nodeMap.entrySet()) { @@ -308,6 +381,8 @@ public class FlowModelBuilder { return NodeTypeEnum.BRANCH; case "loop": return NodeTypeEnum.LOOP; + case "subflow": + return NodeTypeEnum.SUB_FLOW; case "stoploop": return NodeTypeEnum.STOP_LOOP; case "sleep": 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 2522d41..790f7c7 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 @@ -150,6 +150,14 @@ public class TaskInstHolder { .build()); } + /** Records an internal sub-flow node failure without closing the owning task. */ + public void markNodeFailed(String instId, String taskId, String itemId, + String nodeId, String nodeType, String operate, + String action, String params, String message, List iterations) { + nodeInstService.logFailed(instId, taskId, itemId, nodeId, nodeType, + operate, action, params, message, iterations); + } + public void syncStatus(String instId, TaskStatusEnum status) { TaskContext ctx = taskContextManager.get(instId); if (ctx != null) { 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 a2de82c..e12a4a8 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 @@ -115,6 +115,28 @@ public class TaskThreadRegistry { return nodeContextMap; } + /** + * Removes interrupted inner nodes before resume. The owning sub-flow node remains + * registered and is the only restart point, preventing parent and child duplication. + */ + public List discardNestedSubFlowContexts(String instId) { + List removed = new ArrayList<>(); + for (Map.Entry entry : nodeContextMap.entrySet()) { + TaskNodeExecuteContext context = entry.getValue(); + String ownerNodeId = context.getRootMessage() == null + ? null : context.getRootMessage().getSubFlowOwnerNodeId(); + if (!entry.getKey().startsWith(instId + "_") || ownerNodeId == null || ownerNodeId.isBlank()) { + continue; + } + if (nodeContextMap.remove(entry.getKey(), context)) removed.add(context); + } + if (!removed.isEmpty()) { + log.info("恢复前清理子流程内部上下文: instId={}, nodeIds={}", instId, + removed.stream().map(value -> value.getNode().getNodeId()).collect(Collectors.toList())); + } + return removed; + } + /** * 注册 loop 节点 latch(每次只允许一个 latch) */ 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 b1a051b..cd11974 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 @@ -220,8 +220,22 @@ public class FlowControlService { } } + List suspendedNodes = taskThreadRegistry.getPendingNodes(instId); + if (asyncMediaAnalysisCoordinator != null) { + suspendedNodes.stream() + .flatMap(node -> java.util.stream.Stream.of( + node.getActiveSubFlowExecutionId(), + node.getRootMessage() == null ? null : node.getRootMessage().getSubFlowExecutionId())) + .filter(StrUtil::isNotBlank) + .distinct() + .forEach(scopeId -> asyncMediaAnalysisCoordinator.cancelScope(instId, scopeId)); + } + List discardedSubFlowNodes = + taskThreadRegistry.discardNestedSubFlowContexts(instId); + for (TaskNodeExecuteContext discarded : discardedSubFlowNodes) { + nodeInstService.delete(instId, discarded.getNode().getNodeId(), discarded.getIterations()); + } ctx.setPaused(false); - List pendingNodes = taskThreadRegistry.getPendingNodes(instId); List pendingLoops = pendingNodes.stream() .filter(node -> node.getNode().getNodeType().equals(NodeTypeEnum.LOOP)) diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandler.java index 9cf3b81..ab24237 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandler.java @@ -34,7 +34,7 @@ public class FlowEndNodeHandler implements FlowNodeTypeHandler { String taskId = message.getTaskId(); // Earlier items must not wait for media analysis; otherwise the next item // cannot start. The final END node is the task-level completion barrier. - if (message.getPendingItemCount() == 0) { + if (!message.isNestedFlow() && message.getPendingItemCount() == 0) { TaskNodeExecuteResult asyncResult = asyncMediaAnalysisCoordinator.awaitAll(instId); if (!asyncResult.isSuccess()) { return asyncResult; diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowSubFlowNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowSubFlowNodeHandler.java new file mode 100644 index 0000000..25d0c9f --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowSubFlowNodeHandler.java @@ -0,0 +1,100 @@ +package com.cmvr.test.flow.runtime.dispatcher; + +import cn.hutool.core.bean.BeanUtil; +import cn.hutool.core.util.StrUtil; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; +import com.cmvr.common.utils.spring.SpringUtils; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +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.flow.runtime.engine.FlowItemExecutor; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; +import com.cmvr.test.flow.runtime.subflow.FlowSubFlowDefinitionService; +import com.cmvr.test.flow.runtime.subflow.SubFlowExecutionDefinition; +import com.cmvr.test.flow.runtime.subflow.SubFlowExecutionOutcome; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Component; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.UUID; + +@Component("SUB_FLOW") +@RequiredArgsConstructor +public class FlowSubFlowNodeHandler implements FlowNodeTypeHandler { + + private static final int MAX_NESTED_DEPTH = 8; + + private final FlowSubFlowDefinitionService definitionService; + private final TaskInstHolder taskInstHolder; + private final TaskThreadRegistry taskThreadRegistry; + + @Override + public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) { + FlowNodeWrapper node = message.getGraph().getNode(message.getNodeId()); + JSONObject selection = node == null || node.getRawProperties() == null + ? null : node.getRawProperties().getJSONObject("subFlow"); + String childItemId = selection == null ? null : selection.getString("itemId"); + String versionId = selection == null ? null : selection.getString("versionId"); + if (StrUtil.isBlank(childItemId) || StrUtil.isBlank(versionId)) { + throw new GlobalException("子流程节点未选择发布工作流"); + } + definitionService.validateInvocation(message.getItemId(), node.getRawProperties()); + + List stack = message.getSubFlowVersionStack() == null + ? new ArrayList<>() : new ArrayList<>(message.getSubFlowVersionStack()); + if (stack.size() >= MAX_NESTED_DEPTH) { + throw new GlobalException("子流程嵌套层级不能超过 " + MAX_NESTED_DEPTH + " 层"); + } + if (stack.contains(versionId)) { + throw new GlobalException("检测到子流程发布版本循环引用"); + } + stack.add(versionId); + + String rootNodeId = StrUtil.blankToDefault(message.getSubFlowRootNodeId(), message.getNodeId()); + SubFlowExecutionDefinition definition = definitionService.loadExecution( + message.getItemId(), childItemId, versionId, rootNodeId); + String executionId = UUID.randomUUID().toString().replace("-", ""); + SubFlowExecutionOutcome outcome = new SubFlowExecutionOutcome(); + var parentExecutionContext = taskThreadRegistry.getNodeContext( + message.getInstId(), message.getItemId(), message.getNodeId()); + if (parentExecutionContext != null) { + parentExecutionContext.setActiveSubFlowExecutionId(executionId); + } + TaskNodeExecuteMessage childMessage = new TaskNodeExecuteMessage(); + BeanUtil.copyProperties(message, childMessage); + childMessage.setNestedFlow(true); + childMessage.setLocalRunParams(message.getInputParams() == null + ? new JSONObject() : new JSONObject(message.getInputParams())); + childMessage.setSubFlowVersionStack(stack); + childMessage.setSubFlowOwnerNodeId(message.getNodeId()); + childMessage.setSubFlowRootNodeId(rootNodeId); + childMessage.setSubFlowExecutionId(executionId); + childMessage.setSubFlowOutcome(outcome); + + CountDownLatch completed = new CountDownLatch(1); + FlowItemExecutor executor = SpringUtils.getBean(FlowItemExecutor.class); + executor.executeNestedGraph(definition.graph(), childMessage, + completed::countDown, message.getIterations()); + try { + completed.await(); + } catch (InterruptedException error) { + Thread.currentThread().interrupt(); + TaskContext context = taskInstHolder.getContext(message.getInstId()); + if (context == null || context.isPaused() || context.isStopped()) return null; + throw new GlobalException("等待子流程执行被中断", error); + } + + TaskContext context = taskInstHolder.getContext(message.getInstId()); + if (context == null || context.isStopped() || context.isPaused()) return null; + if (outcome.isFailed()) { + return TaskNodeExecuteResult.failure("子流程[" + node.getNodeName() + "]执行失败:" + outcome.getFailure()); + } + JSONObject output = context.getNodeOutput(definition.endNodeId(), message.getIterations()); + return TaskNodeExecuteResult.success(output == null ? new JSONObject() : output); + } +} 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 9ce9225..95dea34 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 @@ -81,18 +81,19 @@ public class FlowItemExecutor { // 主图中开始节点的 inputParams List startInputDefs = subGraph.getStartInputParamsDefs(); - // 子图的 start 节点 - FlowNodeWrapper subStartNode = subGraph.getInDegreeZeroNode(); - String subStartNodeId = subStartNode.getNodeId(); + // A loop body can contain multiple independent entry lines. The first + // entry ID is only resume metadata; the scheduler starts every entry. + List entryNodes = subGraph.getInDegreeZeroNodes(); + if (entryNodes.isEmpty()) { + throw new GlobalException("循环子图没有可执行入口"); + } + String subStartNodeId = entryNodes.get(0).getNodeId(); NodeExecutor executor = node -> executeNode( subGraph, node, subStartNodeId, startInputDefs, rootMessage, onFinished, iterations ); - // 执行子图起始节点 - executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs, - rootMessage, onFinished, iterations); log.info("当前执行-子图"); // 启动子图调度 TaskContext context = taskInstHolder.getContext(rootMessage.getInstId()); @@ -103,6 +104,26 @@ public class FlowItemExecutor { flowTaskScheduler.subStart(subGraph, context, executor, onFinished, iterations); } + /** + * Executes a complete published workflow inside a sub-flow node. Unlike a + * loop body, a published workflow has an explicit START node and must enter + * through it so that input parameters and the full successor chain run. + */ + public void executeNestedGraph(FlowGraph graph, + TaskNodeExecuteMessage rootMessage, + Runnable onFinished, + List iterations) { + String startNodeId = graph.findStartNodeId(); + List startInputDefs = graph.getStartInputParamsDefs(); + TaskContext context = taskInstHolder.getContext(rootMessage.getInstId()); + if (context == null || context.isStopped()) { + notifyFinished(onFinished, rootMessage.getInstId(), startNodeId); + return; + } + executeNode(graph, graph.getStartNode(), startNodeId, startInputDefs, + rootMessage, onFinished, iterations); + } + /** * 执行单个节点(核心调度逻辑) */ @@ -148,11 +169,15 @@ public class FlowItemExecutor { // END 节点 or 子图结束 if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END || graph.isSubGraphEndNode(nodeId)) { - try { - log.info("{} 节点执行,释放线程", node.getNodeType()); - onFinished.run(); - } catch (Exception e) { - log.error("onFinished 回调异常", e); + boolean graphFinished = !graph.isSub() + || flowTaskScheduler.markSubGraphTerminalCompleted(instId, nodeId, graph, iterations); + if (graphFinished) { + try { + log.info("{} 节点执行,释放线程", node.getNodeType()); + onFinished.run(); + } catch (Exception e) { + log.error("onFinished 回调异常", e); + } } } @@ -164,17 +189,30 @@ public class FlowItemExecutor { ); } catch (RuntimeException e) { TaskContext ctx = taskInstHolder.getContext(instId); + boolean nestedFailure = isNestedExecution(rootMessage); 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); + String detail = "节点执行异常: " + e.getMessage(); + if (nestedFailure) { + rootMessage.getSubFlowOutcome().fail(nodeName, detail); + taskInstHolder.markNodeFailed( + instId, rootMessage.getTaskId(), itemId, nodeId, node.getNodeType().name(), + node.getAction() == null ? null : node.getAction().getOperate(), + node.getAction() == null ? null : node.getAction().getAction(), + null, detail, iterations); + } else { + asyncMediaAnalysisCoordinator.cancelAll(instId); + 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, detail, iterations); + } } } finally { notifyFinished(onFinished, instId, nodeId); } + if (nestedFailure || ctx == null || ctx.isPaused() || ctx.isStopped()) return; throw e; } finally { nodeExecutionContext.markExecutionFinished(); @@ -195,6 +233,7 @@ public class FlowItemExecutor { notifyFinished(onFinished, message.getInstId(), message.getNodeId()); } else if (!ctx.isPaused()) { try { + asyncMediaAnalysisCoordinator.cancelAll(message.getInstId()); taskInstHolder.markFailed( message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(), message.getNodeType(), @@ -210,7 +249,20 @@ public class FlowItemExecutor { } try { + if (isNestedExecution(message)) { + String detail = result.getErrorMsg() == null ? "执行失败" : result.getErrorMsg(); + message.getSubFlowOutcome().fail(message.getNodeName(), detail); + taskInstHolder.markNodeFailed( + 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(), + detail, message.getIterations()); + return; + } if (ctx != null && !ctx.isPaused() && !ctx.isStopped()) { + asyncMediaAnalysisCoordinator.cancelAll(message.getInstId()); taskInstHolder.markFailed( message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(), message.getNodeType(), @@ -224,6 +276,10 @@ public class FlowItemExecutor { } } + private boolean isNestedExecution(TaskNodeExecuteMessage message) { + return message != null && message.isNestedFlow() && message.getSubFlowOutcome() != null; + } + private void notifyFinished(Runnable onFinished, String instId, String nodeId) { try { onFinished.run(); @@ -239,7 +295,12 @@ public class FlowItemExecutor { TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(rootMessage, message); message.setNodeId(nodeId); - message.setNodeType(node.getNodeType().name()); + NodeTypeEnum runtimeNodeType = node.getNodeType(); + if (rootMessage.isNestedFlow()) { + if (runtimeNodeType == NodeTypeEnum.START) runtimeNodeType = NodeTypeEnum.SUB_START; + if (runtimeNodeType == NodeTypeEnum.END) runtimeNodeType = NodeTypeEnum.SUB_END; + } + message.setNodeType(runtimeNodeType.name()); message.setAction(node.getAction()); message.setNodeName(nodeName); // message.setInputParams(inputParams); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java index 6ca36b9..5295f40 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java @@ -1,6 +1,7 @@ package com.cmvr.test.flow.runtime.engine; import cn.hutool.core.exceptions.ExceptionUtil; +import com.cmvr.test.enums.NodeTypeEnum; import com.cmvr.test.flow.builder.FlowGraph; import com.cmvr.test.flow.builder.FlowNodeWrapper; import com.cmvr.test.flow.builder.TaskKeyBuilder; @@ -16,6 +17,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Set; +import java.util.stream.Collectors; @Slf4j @Component @@ -58,17 +60,16 @@ public class FlowTaskScheduler { Runnable onFinished, List iterations) { String instId = context.getInstId(); - FlowNodeWrapper startNode = graph.getInDegreeZeroNode(); - List nextNodes = graph.getNextNodes(startNode.getNodeId()); + List entryNodes = graph.getInDegreeZeroNodes(); - if (nextNodes.isEmpty()) { - log.warn("子图 start 节点无出边,流程直接结束:instId={}", instId); + if (entryNodes.isEmpty()) { + log.warn("子图没有可执行入口,流程直接结束:instId={}", instId); onFinished.run(); return; } - for (String nextId : nextNodes) { - scheduleNode(instId, nextId, graph, executorFunc, iterations); + for (FlowNodeWrapper entryNode : entryNodes) { + scheduleNode(instId, entryNode.getNodeId(), graph, executorFunc, iterations); } } @@ -98,6 +99,31 @@ public class FlowTaskScheduler { } } } + + /** + * A loop subgraph may have multiple parallel terminal nodes. Release the + * current iteration only after every remaining terminal path has finished. + */ + public boolean markSubGraphTerminalCompleted(String instId, + String nodeId, + FlowGraph graph, + List iterations) { + Set terminalNodeIds = graph.allNodeIds().stream() + .filter(id -> graph.getOutgoingEdges(id).isEmpty()) + .filter(id -> isExecutableTerminal(graph, id)) + .collect(Collectors.toSet()); + FlowNodeWrapper startNode = graph.getInDegreeZeroNode(); + String graphNodeId = "subgraph-terminal:" + (startNode == null ? "unknown" : startNode.getNodeId()); + String key = TaskKeyBuilder.buildKey(instId, graphNodeId, iterations); + return schedulerCache.markGraphTerminalCompleted(key, nodeId, terminalNodeIds); + } + + private boolean isExecutableTerminal(FlowGraph graph, String nodeId) { + FlowNodeWrapper node = graph.getNode(nodeId); + return node != null + && node.getNodeType() != NodeTypeEnum.START + && node.getNodeType() != NodeTypeEnum.SUB_START; + } /** * 提交节点执行任务(除循环只调度一次) */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java index dce98b0..4619d56 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java @@ -2,6 +2,7 @@ package com.cmvr.test.flow.runtime.engine.support; import cn.hutool.core.collection.CollUtil; import cn.hutool.core.collection.CollectionUtil; +import cn.hutool.core.util.StrUtil; import com.alibaba.fastjson2.JSONArray; import com.alibaba.fastjson2.JSONObject; import com.cmvr.common.exception.GlobalException; @@ -38,7 +39,7 @@ public class FlowNodeParamPreparer { // 3 START 节点:用 runParams 覆盖默认参数 if (node.getNodeType() == NodeTypeEnum.START) { JSONObject merged = new JSONObject(input); - JSONObject runParams = ctx.getRunParams(); + JSONObject runParams = effectiveRunParams(rootMessage, ctx); for (String key : runParams.keySet()) { // runParams 里有值,就覆盖掉定义里的值 merged.put(key, runParams.get(key)); @@ -49,6 +50,12 @@ public class FlowNodeParamPreparer { // 4 LOOP 节点:动态计算 loopCount if (node.getNodeType() == NodeTypeEnum.LOOP) { Object loopNumVal = input.get("loopNum"); + if (isEmptyLoopValue(loopNumVal)) { + loopNumVal = resolveLegacyLoopValue(graph, rootMessage, ctx); + } + if (isEmptyLoopValue(loopNumVal)) { + throw new GlobalException("循环次数未配置,请设置固定次数或引用开始参数"); + } input.put("loopArray", loopNumVal); int loopCount = FlowLoopValueResolver.count(loopNumVal); @@ -59,6 +66,30 @@ public class FlowNodeParamPreparer { return input; } + private Object resolveLegacyLoopValue(FlowGraph graph, TaskNodeExecuteMessage message, TaskContext ctx) { + if (ctx == null) { + return null; + } + JSONObject runParams = effectiveRunParams(message, ctx); + for (String preferredName : List.of("loopNum", "count")) { + Object value = runParams.get(preferredName); + if (!isEmptyLoopValue(value)) { + return value; + } + } + List startDefs = graph.getStartInputParamsDefs(); + if (startDefs != null && startDefs.size() == 1) { + FlowParamDef onlyParam = startDefs.get(0); + Object value = runParams.get(onlyParam.getName()); + return isEmptyLoopValue(value) ? onlyParam.getInput() : value; + } + return null; + } + + private boolean isEmptyLoopValue(Object value) { + return value == null || (value instanceof String && StrUtil.isBlank((String) value)); + } + private JSONObject getInputParams(FlowGraph graph, List paramDefList, TaskNodeExecuteMessage rootMessage) { JSONObject input = new JSONObject(); TaskContext ctx = taskInstHolder.getContext(rootMessage.getInstId()); @@ -80,11 +111,13 @@ public class FlowNodeParamPreparer { // 引用的参数是上游的输入还是输出 // 2.1 引用上游输入参数 if ("input".equalsIgnoreCase(param.getQuoteType())) { - source = ctx.getRunParams(); - JSONObject def = toJson(graph.getNodeParamsDefs().get(refNodeId)); - for (String key : def.keySet()) { - source.put(key, def.get(key)); - } + FlowNodeWrapper referencedNode = graph.getNode(refNodeId); + List referencedDefs = referencedNode != null + && referencedNode.getNodeType() == NodeTypeEnum.START + ? graph.getStartInputParamsDefs() + : graph.getNodeParamsDefs().get(refNodeId); + source = toJson(referencedDefs); + source.putAll(effectiveRunParams(rootMessage, ctx)); } else if ("output".equalsIgnoreCase(param.getQuoteType())) { // 2.2 引用上游的 output 参数(运行时实际执行结果) // 用 nodeId + iterations 作为 key,避免覆盖 @@ -108,6 +141,13 @@ public class FlowNodeParamPreparer { return input; } + private JSONObject effectiveRunParams(TaskNodeExecuteMessage message, TaskContext ctx) { + if (message.getLocalRunParams() != null) { + return message.getLocalRunParams(); + } + return ctx == null || ctx.getRunParams() == null ? new JSONObject() : ctx.getRunParams(); + } + /** * 将参数定义列表中的 input 值构建为 JSONObject diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java index e86e999..bab351c 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java @@ -13,6 +13,8 @@ public class FlowSchedulerCache { // instId_nodeId 或 instId_nodeId_loopIteration private final ConcurrentMap scheduled = new ConcurrentHashMap<>(); private final ConcurrentMap> completedPreMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> completedTerminalMap = new ConcurrentHashMap<>(); + private final ConcurrentMap completedGraphMap = new ConcurrentHashMap<>(); public boolean trySchedule(String key) { AtomicBoolean isScheduled = scheduled.computeIfAbsent(key, k -> new AtomicBoolean(false)); @@ -23,9 +25,22 @@ public class FlowSchedulerCache { return completedPreMap.computeIfAbsent(key, k -> ConcurrentHashMap.newKeySet()); } + public boolean markGraphTerminalCompleted(String key, String nodeId, Set terminalNodeIds) { + Set completedTerminals = completedTerminalMap.computeIfAbsent( + key, ignored -> ConcurrentHashMap.newKeySet()); + completedTerminals.add(nodeId); + if (!completedTerminals.containsAll(terminalNodeIds)) { + return false; + } + AtomicBoolean graphCompleted = completedGraphMap.computeIfAbsent(key, ignored -> new AtomicBoolean(false)); + return graphCompleted.compareAndSet(false, true); + } + public void clearContextCache(String instId) { String prefix = instId + ":"; scheduled.keySet().removeIf(k -> k.startsWith(prefix)); completedPreMap.keySet().removeIf(k -> k.startsWith(prefix)); + completedTerminalMap.keySet().removeIf(k -> k.startsWith(prefix)); + completedGraphMap.keySet().removeIf(k -> k.startsWith(prefix)); } } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java index cdb5d88..3187ca0 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptor.java @@ -10,6 +10,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.core.annotation.Order; import org.springframework.stereotype.Component; +import java.util.ArrayList; import java.util.function.Function; @Slf4j @@ -47,6 +48,8 @@ public class FlowBeforeInterceptor extends AbstractFlowMsgPreInterceptor{ .nodeName(message.getNodeName()) .nodeType(message.getNodeType()) .action(message.getAction()) + .iterations(message.getIterations() == null + ? new ArrayList<>() : new ArrayList<>(message.getIterations())) .build(); eventPublisher.publishEvent(event); 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 fa13acc..dbeefbb 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 @@ -58,7 +58,9 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor { taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), result.getOutputParams() == null ? null : result.getOutputParams().toJSONString(), - StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations()); + StrUtil.format("[{}] 执行成功", + StrUtil.blankToDefault(message.getNodeName(), String.valueOf(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/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java index c60127a..b63c9ba 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java @@ -28,6 +28,7 @@ public class TaskNodeExecuteContext { private int loopIteration; private Object loopArray; private List iterations = new ArrayList<>(); + private String activeSubFlowExecutionId; public void markExecutionFinished() { executionFinished.countDown(); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java index 0c5bd59..ff4f62b 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java @@ -3,6 +3,7 @@ package com.cmvr.test.flow.runtime.message; import com.alibaba.fastjson2.JSONObject; import com.cmvr.test.enums.ActionEnum; import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.runtime.subflow.SubFlowExecutionOutcome; import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO; import lombok.Data; @@ -106,6 +107,27 @@ public class TaskNodeExecuteMessage { */ private JSONObject inputParams; + /** Parameters scoped to a nested published workflow invocation. */ + private JSONObject localRunParams; + + /** True while executing a published workflow inside a sub-flow node. */ + private boolean nestedFlow; + + /** Immutable published versions currently active in the nested call chain. */ + private List subFlowVersionStack = new ArrayList<>(); + + /** Immediate sub-flow node that owns this nested execution context. */ + private String subFlowOwnerNodeId; + + /** Top-level canvas sub-flow node used to group nested trial-run logs. */ + private String subFlowRootNodeId; + + /** Unique invocation scope used to wait for nested asynchronous analysis jobs. */ + private String subFlowExecutionId; + + /** Shared outcome for all nodes in the current published sub-flow invocation. */ + private transient SubFlowExecutionOutcome subFlowOutcome; + /** * 上游输出参数(来自上一个节点的输出) */ diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinator.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinator.java index cfc7917..d3d6de2 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinator.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinator.java @@ -24,7 +24,9 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CancellationException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; +import java.util.Set; import java.util.function.Supplier; @Slf4j @@ -69,7 +71,15 @@ public class AsyncMediaAnalysisCoordinator { completeNode(message, moduleCode, result, error, group, nodeKey); return null; }); - group.jobs.put(nodeKey, completion); + String jobKey = nodeKey + "#" + group.sequence.incrementAndGet(); + group.jobs.put(jobKey, completion); + group.nodeJobs.put(nodeKey, completion); + if (StrUtil.isNotBlank(message.getSubFlowExecutionId())) { + group.scopeJobs.computeIfAbsent(message.getSubFlowExecutionId(), ignored -> new ConcurrentHashMap<>()) + .put(jobKey, completion); + group.scopeNodeKeys.computeIfAbsent(message.getSubFlowExecutionId(), ignored -> ConcurrentHashMap.newKeySet()) + .add(nodeKey); + } log.info("异步媒体分析已提交:instId={}, nodeId={}, requestId={}", message.getInstId(), message.getNodeId(), requestId); return TaskNodeExecuteResult.deferred(accepted); @@ -80,7 +90,7 @@ public class AsyncMediaAnalysisCoordinator { ExecutionGroup group = groups.get(instId); if (group == null) return; String key = nodeKey(itemId, itemOrder, itemOccurrence, nodeId, safeIterations(iterations)); - CompletableFuture job = group.jobs.get(key); + CompletableFuture job = group.nodeJobs.get(key); if (job != null) await(job); String error = group.nodeFailures.get(key); if (error != null) throw new GlobalException(error); @@ -97,13 +107,34 @@ public class AsyncMediaAnalysisCoordinator { if (group.jobs.values().stream().allMatch(CompletableFuture::isDone)) break; } completed = true; - String error = group.failure.get(); + String error = group.nodeFailures.values().stream().findFirst().orElse(null); return error == null ? TaskNodeExecuteResult.success() : TaskNodeExecuteResult.failure(error); } finally { if (completed) groups.remove(instId, group); } } + /** Waits only for asynchronous analysis submitted by one sub-flow invocation. */ + public TaskNodeExecuteResult awaitScope(String instId, String executionId) { + if (StrUtil.isBlank(executionId)) return TaskNodeExecuteResult.success(); + ExecutionGroup group = groups.get(instId); + if (group == null) return TaskNodeExecuteResult.success(); + Map> scopedJobs = group.scopeJobs.get(executionId); + if (scopedJobs == null) return TaskNodeExecuteResult.success(); + try { + while (true) { + CompletableFuture[] snapshot = scopedJobs.values().toArray(new CompletableFuture[0]); + await(CompletableFuture.allOf(snapshot)); + if (scopedJobs.values().stream().allMatch(CompletableFuture::isDone)) break; + } + String error = group.scopeFailures.get(executionId); + return error == null ? TaskNodeExecuteResult.success() : TaskNodeExecuteResult.failure(error); + } finally { + group.scopeJobs.remove(executionId, scopedJobs); + group.scopeFailures.remove(executionId); + } + } + public void cancelAll(String instId) { ExecutionGroup group = groups.remove(instId); if (group == null) return; @@ -112,10 +143,34 @@ public class AsyncMediaAnalysisCoordinator { log.info("已取消任务的异步媒体分析登记:instId={}, count={}", instId, group.jobs.size()); } + /** Cancels and detaches jobs from one interrupted sub-flow attempt before resume. */ + public void cancelScope(String instId, String executionId) { + if (StrUtil.isBlank(executionId)) return; + ExecutionGroup group = groups.get(instId); + if (group == null) return; + group.cancelledScopes.add(executionId); + Map> scopedJobs = group.scopeJobs.remove(executionId); + group.scopeFailures.remove(executionId); + Set scopedNodeKeys = group.scopeNodeKeys.remove(executionId); + if (scopedNodeKeys != null) scopedNodeKeys.forEach(group.nodeFailures::remove); + if (scopedJobs == null) return; + scopedJobs.forEach((jobKey, job) -> { + group.jobs.remove(jobKey, job); + group.nodeJobs.entrySet().removeIf(entry -> entry.getValue() == job); + job.cancel(true); + }); + log.info("已取消子流程异步分析作用域: instId={}, executionId={}, count={}", + instId, executionId, scopedJobs.size()); + } + private void completeNode(TaskNodeExecuteMessage message, String moduleCode, TaskNodeExecuteResult result, Throwable error, ExecutionGroup group, String nodeKey) { try { + if (StrUtil.isNotBlank(message.getSubFlowExecutionId()) + && group.cancelledScopes.contains(message.getSubFlowExecutionId())) { + return; + } TaskContext currentContext = taskInstHolder.getContext(message.getInstId()); if (group.cancelled || currentContext == null || currentContext.isStopped()) { log.info("任务已终止,忽略异步媒体分析结果:instId={}, nodeId={}", @@ -131,6 +186,7 @@ public class AsyncMediaAnalysisCoordinator { StrUtil.blankToDefault(ExceptionUtil.getMessage(cause), cause.toString()), 0, 480); group.failure.compareAndSet(null, detail); group.nodeFailures.put(nodeKey, detail); + recordScopeFailure(group, message, detail); nodeInstService.logFailed(message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(), message.getNodeType(), message.getAction() == null ? null : message.getAction().getOperate(), @@ -145,13 +201,16 @@ public class AsyncMediaAnalysisCoordinator { taskInstHolder.setNodeOutParams(message.getInstId(), message.getNodeId(), safeIterations(message), output); taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), output.toJSONString(), - "异步分析执行成功", safeIterations(message)); + StrUtil.format("[{}] 执行成功", + StrUtil.blankToDefault(message.getNodeName(), "异步分析")), + safeIterations(message)); publish(message, moduleCode, FlowExecutionEvent.EventType.NODE_COMPLETED, output, null); } catch (Throwable finalizeError) { String detail = "异步分析结果处理失败:" + StrUtil.sub( StrUtil.blankToDefault(ExceptionUtil.getMessage(finalizeError), finalizeError.toString()), 0, 480); group.failure.compareAndSet(null, detail); group.nodeFailures.put(nodeKey, detail); + recordScopeFailure(group, message, detail); log.error("异步媒体分析结果处理异常:instId={}, nodeId={}", message.getInstId(), message.getNodeId(), finalizeError); } @@ -180,6 +239,12 @@ public class AsyncMediaAnalysisCoordinator { .build()); } + private void recordScopeFailure(ExecutionGroup group, TaskNodeExecuteMessage message, String detail) { + if (StrUtil.isNotBlank(message.getSubFlowExecutionId())) { + group.scopeFailures.putIfAbsent(message.getSubFlowExecutionId(), detail); + } + } + private void runDetached(TaskNodeExecuteMessage message, Supplier analysis) { try { analysis.get(); @@ -220,7 +285,13 @@ public class AsyncMediaAnalysisCoordinator { private static final class ExecutionGroup { private final Map> jobs = new ConcurrentHashMap<>(); + private final Map> nodeJobs = new ConcurrentHashMap<>(); private final Map nodeFailures = new ConcurrentHashMap<>(); + private final Map>> scopeJobs = new ConcurrentHashMap<>(); + private final Map scopeFailures = new ConcurrentHashMap<>(); + private final Map> scopeNodeKeys = new ConcurrentHashMap<>(); + private final Set cancelledScopes = ConcurrentHashMap.newKeySet(); + private final AtomicLong sequence = new AtomicLong(); private final AtomicReference failure = new AtomicReference<>(); private volatile boolean cancelled; } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionService.java new file mode 100644 index 0000000..1f8779b --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionService.java @@ -0,0 +1,480 @@ +package com.cmvr.test.flow.runtime.subflow; + +import cn.hutool.core.util.StrUtil; +import cn.hutool.crypto.digest.DigestUtil; +import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONArray; +import com.alibaba.fastjson2.JSONObject; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.cmvr.common.exception.GlobalException; +import com.cmvr.test.enums.NodeTypeEnum; +import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.builder.FlowModelBuilder; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +import com.cmvr.test.mapper.TeDetectionItemMapper; +import com.cmvr.test.mapper.TeDetectionItemVersionMapper; +import com.cmvr.test.model.domain.TeDetectionItem; +import com.cmvr.test.model.domain.TeDetectionItemVersion; +import com.cmvr.test.model.vo.TeSubFlowOptionVO; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.Comparator; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + +@Service +@RequiredArgsConstructor +public class FlowSubFlowDefinitionService { + + public static final String PROJECT_AIMA = "AIMA"; + public static final String PROJECT_CHANGAN = "CHANGAN"; + private static final int MAX_SUB_FLOW_DEPTH = 8; + + private final TeDetectionItemMapper itemMapper; + private final TeDetectionItemVersionMapper versionMapper; + + public List listOptions(String currentItemId) { + TeDetectionItem current = requireItem(currentItemId); + Set aimaItemIds = new HashSet<>(itemMapper.selectAimaDetectionItemIds()); + String projectCode = projectCode(current.getId(), aimaItemIds); + List items = itemMapper.selectList(new LambdaQueryWrapper() + .isNotNull(TeDetectionItem::getPublishedVersionId) + .ne(TeDetectionItem::getPublishedVersionId, "") + .ne(TeDetectionItem::getId, currentItemId) + .eq(TeDetectionItem::getStatus, "0")); + List result = new ArrayList<>(); + for (TeDetectionItem item : items) { + if (!projectCode.equals(projectCode(item.getId(), aimaItemIds))) continue; + TeDetectionItemVersion version = versionMapper.selectById(item.getPublishedVersionId()); + if (version == null || createsCircularReference(current.getId(), item.getId(), version)) continue; + result.add(toContract(item, version, projectCode)); + } + result.sort(Comparator.comparing(TeSubFlowOptionVO::getDetectName, + Comparator.nullsLast(String::compareToIgnoreCase))); + return result; + } + + public void validateForPublish(TeDetectionItem parent, String flowData) { + JSONObject flow = parseFlow(flowData); + JSONArray nodes = flow.getJSONArray("nodes"); + if (nodes == null) return; + validateNoCircularReference(parent, flow); + validateEndOutputs(flow); + for (int i = 0; i < nodes.size(); i++) { + JSONObject node = nodes.getJSONObject(i); + JSONObject properties = node.getJSONObject("properties"); + if (!isSubFlowNode(node, properties)) continue; + JSONObject selection = properties == null ? null : properties.getJSONObject("subFlow"); + String childItemId = selection == null ? null : selection.getString("itemId"); + String versionId = selection == null ? null : selection.getString("versionId"); + if (StrUtil.isBlank(childItemId) || StrUtil.isBlank(versionId)) { + throw new GlobalException("子流程节点[" + nodeName(properties) + "]未选择发布工作流"); + } + if (parent.getId().equals(childItemId)) { + throw new GlobalException("子流程节点不能引用当前工作流自身"); + } + TeDetectionItem child = requireSameProject(parent.getId(), childItemId); + if (!versionId.equals(child.getPublishedVersionId())) { + throw new GlobalException("子流程[" + child.getDetectName() + "]已发布新版本,请重新同步后再发布"); + } + TeDetectionItemVersion version = requireVersion(childItemId, versionId); + TeSubFlowOptionVO contract = toContract(child, version, projectCode(parent.getId())); + validateContract(nodeName(properties), properties, contract); + validateReferences(flow, node, contract.getInputParams(), "子流程节点[" + nodeName(properties) + "]的入参"); + validateResources(flow, contract, child.getDetectName()); + } + } + + /** + * Published versions are immutable, but a newly published parent may still close an + * item-level chain such as A -> B -> A. Validate the complete selected-version tree + * before publishing so the error is reported while editing instead of at runtime. + */ + private void validateNoCircularReference(TeDetectionItem parent, JSONObject flow) { + List path = new ArrayList<>(); + path.add(parent.getId()); + CircularReference circular = findCircularReference(flow, path, 0); + if (circular != null) { + if (circular.depthExceeded()) { + throw new GlobalException("子流程嵌套层级不能超过 " + MAX_SUB_FLOW_DEPTH + " 层:" + + formatItemPath(circular.itemIds())); + } + throw new GlobalException("检测到子流程循环引用:" + formatItemPath(circular.itemIds())); + } + } + + private boolean createsCircularReference(String currentItemId, String candidateItemId, + TeDetectionItemVersion candidateVersion) { + List path = new ArrayList<>(); + path.add(currentItemId); + if (path.contains(candidateItemId)) return true; + path.add(candidateItemId); + return findCircularReference(parseFlow(candidateVersion.getFlowData()), path, 1) != null; + } + + private CircularReference findCircularReference(JSONObject flow, List path, int depth) { + JSONArray nodes = flow.getJSONArray("nodes"); + if (nodes == null) return null; + for (int i = 0; i < nodes.size(); i++) { + JSONObject node = nodes.getJSONObject(i); + JSONObject properties = node.getJSONObject("properties"); + if (!isSubFlowNode(node, properties)) continue; + JSONObject selection = properties == null ? null : properties.getJSONObject("subFlow"); + String childItemId = selection == null ? null : selection.getString("itemId"); + String versionId = selection == null ? null : selection.getString("versionId"); + if (StrUtil.isBlank(childItemId) || StrUtil.isBlank(versionId)) continue; + + int repeatedAt = path.indexOf(childItemId); + if (repeatedAt >= 0) { + List cycle = new ArrayList<>(path.subList(repeatedAt, path.size())); + cycle.add(childItemId); + return new CircularReference(cycle, false); + } + if (depth >= MAX_SUB_FLOW_DEPTH) { + List excessivePath = new ArrayList<>(path); + excessivePath.add(childItemId); + return new CircularReference(excessivePath, true); + } + + TeDetectionItemVersion version = versionMapper.selectById(versionId); + if (version == null || !childItemId.equals(version.getDetectionItemId()) + || StrUtil.isBlank(version.getFlowData())) { + continue; + } + path.add(childItemId); + CircularReference nested = findCircularReference(parseFlow(version.getFlowData()), path, depth + 1); + path.remove(path.size() - 1); + if (nested != null) return nested; + } + return null; + } + + private String formatItemPath(List itemIds) { + List names = new ArrayList<>(); + for (String itemId : itemIds) { + TeDetectionItem item = itemMapper.selectById(itemId); + names.add(item == null ? itemId : StrUtil.blankToDefault(item.getDetectName(), itemId)); + } + return String.join(" -> ", names); + } + + private record CircularReference(List itemIds, boolean depthExceeded) { + } + + public SubFlowExecutionDefinition loadExecution(String parentItemId, String childItemId, + String versionId, String parentNodeId) { + if (StrUtil.isBlank(parentNodeId)) throw new GlobalException("子流程节点ID不能为空"); + requireSameProject(parentItemId, childItemId); + TeDetectionItemVersion version = requireVersion(childItemId, versionId); + JSONObject namespaced = namespace(parseFlow(version.getFlowData()), parentNodeId); + FlowGraph graph = FlowModelBuilder.buildExecutableGraph(namespaced.toJSONString()); + String endNodeId = graph.getNodeMap().values().stream() + .filter(node -> node.getNodeType() == NodeTypeEnum.END) + .map(FlowNodeWrapper::getNodeId) + .findFirst() + .orElseThrow(() -> new GlobalException("子流程缺少结束节点")); + return new SubFlowExecutionDefinition(graph, endNodeId, versionId); + } + + public void validateInvocation(String parentItemId, JSONObject properties) { + JSONObject selection = properties == null ? null : properties.getJSONObject("subFlow"); + String childItemId = selection == null ? null : selection.getString("itemId"); + String versionId = selection == null ? null : selection.getString("versionId"); + TeDetectionItem child = requireSameProject(parentItemId, childItemId); + TeDetectionItemVersion version = requireVersion(childItemId, versionId); + validateContract(nodeName(properties), properties, + toContract(child, version, projectCode(parentItemId))); + } + + private TeDetectionItem requireSameProject(String parentItemId, String childItemId) { + TeDetectionItem parent = requireItem(parentItemId); + TeDetectionItem child = requireItem(childItemId); + Set aimaItemIds = new HashSet<>(itemMapper.selectAimaDetectionItemIds()); + if (!projectCode(parent.getId(), aimaItemIds).equals(projectCode(child.getId(), aimaItemIds))) { + throw new GlobalException("子流程与当前工作流不属于同一项目"); + } + return child; + } + + private TeDetectionItem requireItem(String itemId) { + TeDetectionItem item = StrUtil.isBlank(itemId) ? null : itemMapper.selectById(itemId); + if (item == null) throw new GlobalException("工作流不存在: " + itemId); + return item; + } + + private TeDetectionItemVersion requireVersion(String itemId, String versionId) { + TeDetectionItemVersion version = StrUtil.isBlank(versionId) ? null : versionMapper.selectById(versionId); + if (version == null || !itemId.equals(version.getDetectionItemId()) || StrUtil.isBlank(version.getFlowData())) { + throw new GlobalException("子流程发布版本不存在或已失效"); + } + return version; + } + + private TeSubFlowOptionVO toContract(TeDetectionItem item, TeDetectionItemVersion version, String projectCode) { + JSONObject flow = parseFlow(version.getFlowData()); + JSONObject start = findNodeProperties(flow, "start"); + JSONObject end = findNodeProperties(flow, "end"); + if (start == null || end == null) throw new GlobalException("发布工作流缺少开始或结束节点: " + item.getDetectName()); + TeSubFlowOptionVO result = new TeSubFlowOptionVO(); + result.setItemId(item.getId()); + result.setDetectName(item.getDetectName()); + result.setProjectCode(projectCode); + result.setVersionId(version.getId()); + result.setVersionNo(version.getVersionNo()); + result.setInputParams(copyArray(start.getJSONArray("inputParams"))); + result.setOutputParams(copyArray(end.getJSONArray("outputParams"))); + result.setResourceRoles(copyArray(start.getJSONArray("resourceRoles"))); + return result; + } + + private void validateContract(String nodeName, JSONObject properties, TeSubFlowOptionVO contract) { + JSONArray params = properties.getJSONArray("nodeParams"); + Map paramsByName = byName(params); + for (Object value : contract.getInputParams()) { + JSONObject input = (JSONObject) value; + JSONObject param = paramsByName.get(input.getString("name")); + JSONArray quote = param == null ? null : param.getJSONArray("quote"); + if (param == null || !"quote".equals(param.getString("type")) || quote == null || quote.size() < 3) { + throw new GlobalException("子流程节点[" + nodeName + "]的入参[" + input.getString("name") + "]必须引用上游参数"); + } + } + if (!contractSignature(contract.getOutputParams()).equals( + contractSignature(properties.getJSONArray("outputParams")))) { + throw new GlobalException("子流程节点[" + nodeName + "]的输出定义已过期,请重新同步"); + } + } + + private void validateEndOutputs(JSONObject flow) { + JSONArray nodes = flow.getJSONArray("nodes"); + for (Object value : nodes) { + JSONObject node = (JSONObject) value; + if (!"end".equalsIgnoreCase(node.getString("type"))) continue; + JSONObject properties = node.getJSONObject("properties"); + validateReferences(flow, node, + properties == null ? null : properties.getJSONArray("outputParams"), + "结束节点输出"); + } + } + + private void validateReferences(JSONObject flow, JSONObject targetNode, JSONArray contractParams, String label) { + if (contractParams == null || contractParams.isEmpty()) return; + JSONObject targetProperties = targetNode.getJSONObject("properties"); + Map paramsByName = byName(targetProperties == null + ? null : targetProperties.getJSONArray("nodeParams")); + Set upstreamIds = upstreamNodeIds(flow, targetNode.getString("id")); + Map nodesById = nodesById(flow); + for (Object value : contractParams) { + JSONObject contractParam = (JSONObject) value; + String paramName = contractParam.getString("name"); + JSONObject param = paramsByName.get(paramName); + JSONArray quote = param == null ? null : param.getJSONArray("quote"); + if (param == null || !"quote".equals(param.getString("type")) || quote == null || quote.size() < 3) { + throw new GlobalException(label + "[" + paramName + "]必须引用开始节点输入或上游节点输出"); + } + String sourceNodeId = quote.getString(0); + JSONObject sourceNode = nodesById.get(sourceNodeId); + if (!upstreamIds.contains(sourceNodeId) || sourceNode == null) { + throw new GlobalException(label + "[" + paramName + "]引用的节点不是当前节点的上游节点"); + } + boolean startSource = "start".equalsIgnoreCase(sourceNode.getString("type")); + String expectedKind = startSource ? "input" : "output"; + if (!expectedKind.equalsIgnoreCase(quote.getString(1))) { + throw new GlobalException(label + "[" + paramName + "]只能引用开始节点输入或上游节点输出"); + } + JSONObject sourceProperties = sourceNode.getJSONObject("properties"); + JSONArray sourceParams = sourceProperties == null ? null + : sourceProperties.getJSONArray(startSource ? "inputParams" : "outputParams"); + if (!byName(sourceParams).containsKey(quote.getString(2))) { + throw new GlobalException(label + "[" + paramName + "]引用的参数不存在"); + } + } + } + + private Set upstreamNodeIds(JSONObject flow, String targetNodeId) { + Set result = new HashSet<>(); + JSONArray nodes = flow.getJSONArray("nodes"); + for (Object value : nodes) { + JSONObject node = (JSONObject) value; + if ("start".equalsIgnoreCase(node.getString("type"))) result.add(node.getString("id")); + } + Set pending = new HashSet<>(); + pending.add(targetNodeId); + JSONArray edges = flow.getJSONArray("edges"); + while (!pending.isEmpty()) { + String target = pending.iterator().next(); + pending.remove(target); + if (edges == null) continue; + for (Object value : edges) { + JSONObject edge = (JSONObject) value; + if (!target.equals(edge.getString("targetNodeId"))) continue; + String source = edge.getString("sourceNodeId"); + if (targetNodeId.equals(source)) continue; + if (result.add(source)) pending.add(source); + } + } + return result; + } + + private Map nodesById(JSONObject flow) { + Map result = new HashMap<>(); + for (Object value : flow.getJSONArray("nodes")) { + JSONObject node = (JSONObject) value; + result.put(node.getString("id"), node); + } + return result; + } + + private void validateResources(JSONObject parentFlow, TeSubFlowOptionVO contract, String childName) { + JSONObject parentStart = findNodeProperties(parentFlow, "start"); + Map parentRoles = byKey(parentStart == null ? null + : parentStart.getJSONArray("resourceRoles"), "roleKey"); + for (Object value : contract.getResourceRoles()) { + JSONObject childRole = (JSONObject) value; + JSONObject parentRole = parentRoles.get(childRole.getString("roleKey")); + if (parentRole == null) throw new GlobalException("子流程[" + childName + "]所需机器人角色未同步到开始节点"); + Map parentSlots = byKey(parentRole.getJSONArray("deviceSlots"), "slotKey"); + JSONArray childSlots = childRole.getJSONArray("deviceSlots"); + if (childSlots == null) continue; + for (Object slotValue : childSlots) { + JSONObject childSlot = (JSONObject) slotValue; + JSONObject parentSlot = parentSlots.get(childSlot.getString("slotKey")); + if (parentSlot == null || !java.util.Objects.equals( + parentSlot.getInteger("deviceKind"), childSlot.getInteger("deviceKind"))) { + throw new GlobalException("子流程[" + childName + "]所需设备用途未同步到开始节点"); + } + } + } + } + + private JSONObject namespace(JSONObject flow, String rootNodeId) { + final int maxRuntimeNodeIdLength = 32; + final String marker = "sf_"; + String rootPrefix = rootNodeId.substring(0, Math.min(12, rootNodeId.length())); + int hashLength = maxRuntimeNodeIdLength - marker.length() - rootPrefix.length() - 1; + if (hashLength < 12) throw new GlobalException("子流程节点ID过长,无法生成运行日志标识"); + String invocationToken = UUID.randomUUID().toString(); + JSONArray nodes = flow.getJSONArray("nodes"); + Map ids = new HashMap<>(); + for (Object value : nodes) { + JSONObject node = (JSONObject) value; + String oldId = node.getString("id"); + String hash = DigestUtil.sha256Hex(rootNodeId + ":" + invocationToken + ":" + oldId) + .substring(0, hashLength); + ids.put(oldId, marker + rootPrefix + "_" + hash); + } + for (Object value : nodes) { + JSONObject node = (JSONObject) value; + String oldId = node.getString("id"); + node.put("id", ids.get(oldId)); + rewriteIdArray(node.getJSONArray("children"), ids); + JSONObject properties = node.getJSONObject("properties"); + if (properties != null && ids.containsKey(properties.getString("parentId"))) { + properties.put("parentId", ids.get(properties.getString("parentId"))); + } + rewriteReferences(properties, ids); + } + JSONArray edges = flow.getJSONArray("edges"); + if (edges != null) { + for (Object value : edges) { + JSONObject edge = (JSONObject) value; + edge.put("id", "sf_" + UUID.randomUUID().toString().replace("-", "")); + edge.put("sourceNodeId", ids.get(edge.getString("sourceNodeId"))); + edge.put("targetNodeId", ids.get(edge.getString("targetNodeId"))); + } + } + return flow; + } + + private void rewriteReferences(Object value, Map ids) { + if (value instanceof JSONObject object) { + for (String key : new ArrayList<>(object.keySet())) { + Object child = object.get(key); + if (("quote".equals(key) || "nameQuote".equals(key)) && child instanceof JSONArray quote + && !quote.isEmpty() && ids.containsKey(quote.getString(0))) { + quote.set(0, ids.get(quote.getString(0))); + } else { + rewriteReferences(child, ids); + } + } + } else if (value instanceof JSONArray array) { + for (Object child : array) rewriteReferences(child, ids); + } + } + + private void rewriteIdArray(JSONArray values, Map ids) { + if (values == null) return; + for (int i = 0; i < values.size(); i++) { + if (ids.containsKey(values.getString(i))) values.set(i, ids.get(values.getString(i))); + } + } + + private JSONObject parseFlow(String flowData) { + try { + JSONObject flow = JSON.parseObject(flowData); + if (flow == null || flow.getJSONArray("nodes") == null) throw new IllegalArgumentException(); + return flow; + } catch (RuntimeException error) { + throw new GlobalException("子流程发布版本数据格式错误"); + } + } + + private JSONObject findNodeProperties(JSONObject flow, String type) { + JSONArray nodes = flow.getJSONArray("nodes"); + for (Object value : nodes) { + JSONObject node = (JSONObject) value; + if (type.equalsIgnoreCase(node.getString("type"))) return node.getJSONObject("properties"); + } + return null; + } + + private boolean isSubFlowNode(JSONObject node, JSONObject properties) { + return "subFlow".equalsIgnoreCase(node.getString("type")) + || properties != null && "SUB_FLOW".equalsIgnoreCase(properties.getString("action")); + } + + private String nodeName(JSONObject properties) { + return properties == null ? "子流程" : StrUtil.blankToDefault(properties.getString("name"), "子流程"); + } + + private JSONArray copyArray(JSONArray source) { + return source == null ? new JSONArray() : JSON.parseArray(source.toJSONString()); + } + + private Map byName(JSONArray values) { + return byKey(values, "name"); + } + + private Map byKey(JSONArray values, String key) { + Map result = new HashMap<>(); + if (values == null) return result; + for (Object value : values) { + JSONObject object = (JSONObject) value; + result.put(object.getString(key), object); + } + return result; + } + + private List contractSignature(JSONArray values) { + List result = new ArrayList<>(); + if (values == null) return result; + for (Object value : values) { + JSONObject param = (JSONObject) value; + result.add(param.getString("name") + ":" + param.getString("type")); + } + return result; + } + + private String projectCode(String itemId) { + return projectCode(itemId, new HashSet<>(itemMapper.selectAimaDetectionItemIds())); + } + + private String projectCode(String itemId, Set aimaItemIds) { + return aimaItemIds.contains(itemId) ? PROJECT_AIMA : PROJECT_CHANGAN; + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionDefinition.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionDefinition.java new file mode 100644 index 0000000..ddf8b2f --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionDefinition.java @@ -0,0 +1,6 @@ +package com.cmvr.test.flow.runtime.subflow; + +import com.cmvr.test.flow.builder.FlowGraph; + +public record SubFlowExecutionDefinition(FlowGraph graph, String endNodeId, String versionId) { +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcome.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcome.java new file mode 100644 index 0000000..b099364 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcome.java @@ -0,0 +1,25 @@ +package com.cmvr.test.flow.runtime.subflow; + +import cn.hutool.core.util.StrUtil; + +import java.util.concurrent.atomic.AtomicReference; + +/** Carries the first nested execution failure back to the owning sub-flow node. */ +public class SubFlowExecutionOutcome { + + private final AtomicReference failure = new AtomicReference<>(); + + public void fail(String nodeName, String detail) { + String name = StrUtil.blankToDefault(nodeName, "未知节点"); + String reason = StrUtil.blankToDefault(detail, "执行失败"); + failure.compareAndSet(null, "内部节点[" + name + "]" + reason); + } + + public boolean isFailed() { + return failure.get() != null; + } + + public String getFailure() { + return failure.get(); + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/mapper/TeDetectionItemMapper.java b/cmvr-iot-test/src/main/java/com/cmvr/test/mapper/TeDetectionItemMapper.java index e6460c7..1199f76 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/mapper/TeDetectionItemMapper.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/mapper/TeDetectionItemMapper.java @@ -2,6 +2,9 @@ package com.cmvr.test.mapper; import com.cmvr.test.model.domain.TeDetectionItem; import com.github.yulichang.base.MPJBaseMapper; +import org.apache.ibatis.annotations.Select; + +import java.util.List; /** * 检测项配置Mapper接口 @@ -9,4 +12,7 @@ import com.github.yulichang.base.MPJBaseMapper; * @author cmvr-iot */ public interface TeDetectionItemMapper extends MPJBaseMapper { + + @Select("select distinct detect_item_id from aima_test_case where detect_item_id is not null") + List selectAimaDetectionItemIds(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeSubFlowOptionVO.java b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeSubFlowOptionVO.java new file mode 100644 index 0000000..ef98527 --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/model/vo/TeSubFlowOptionVO.java @@ -0,0 +1,16 @@ +package com.cmvr.test.model.vo; + +import com.alibaba.fastjson2.JSONArray; +import lombok.Data; + +@Data +public class TeSubFlowOptionVO { + private String itemId; + private String detectName; + private String projectCode; + private String versionId; + private Integer versionNo; + private JSONArray inputParams; + private JSONArray outputParams; + private JSONArray resourceRoles; +} 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 1dfcc9f..0973913 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 @@ -92,6 +92,7 @@ public class ITeNodeInstServiceImpl extends ServiceImpl FlowModelBuilder.buildExecutableGraph(flow)); + + assertEquals("节点参数存在循环引用:A -> B -> A", error.getMessage()); + } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java index f43ac28..424704f 100644 --- a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/context/TaskThreadRegistryTest.java @@ -1,24 +1,41 @@ package com.cmvr.test.flow.context; +import com.cmvr.test.flow.builder.FlowNodeWrapper; import com.cmvr.test.flow.runtime.message.TaskNodeExecuteContext; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; import org.junit.Test; -import static org.junit.Assert.assertSame; +import java.util.List; + +import static org.junit.Assert.assertEquals; public class TaskThreadRegistryTest { @Test - public void staleExecutionCannotUnregisterResumedNodeContext() { + public void resumeDiscardKeepsOwnerAndRemovesNestedSubFlowNodes() { TaskThreadRegistry registry = new TaskThreadRegistry(); - TaskNodeExecuteContext interrupted = new TaskNodeExecuteContext(); - TaskNodeExecuteContext resumed = new TaskNodeExecuteContext(); - interrupted.setThread(Thread.currentThread()); - resumed.setThread(Thread.currentThread()); + TaskNodeExecuteContext owner = context("sub-flow-node", null); + TaskNodeExecuteContext nested = context("sub-flow-node__subflow__run__child", "sub-flow-node"); + registry.registerNodeContext("inst-1", "item-1", owner.getNode().getNodeId(), owner); + registry.registerNodeContext("inst-1", "item-1", nested.getNode().getNodeId(), nested); - registry.registerNodeContext("inst-1", "item-1", "node-1", interrupted); - registry.registerNodeContext("inst-1", "item-1", "node-1", resumed); - registry.unregisterNodeContext("inst-1", "item-1", "node-1", interrupted); + List removed = registry.discardNestedSubFlowContexts("inst-1"); - assertSame(resumed, registry.getNodeContext("inst-1", "item-1", "node-1")); + assertEquals(1, removed.size()); + assertEquals(nested, removed.get(0)); + assertEquals(List.of(owner), registry.getPendingNodes("inst-1")); + } + + private TaskNodeExecuteContext context(String nodeId, String ownerNodeId) { + FlowNodeWrapper node = new FlowNodeWrapper(); + node.setNodeId(nodeId); + TaskNodeExecuteMessage rootMessage = new TaskNodeExecuteMessage(); + rootMessage.setItemId("item-1"); + rootMessage.setSubFlowOwnerNodeId(ownerNodeId); + TaskNodeExecuteContext context = new TaskNodeExecuteContext(); + context.setNode(node); + context.setRootMessage(rootMessage); + context.setThread(Thread.currentThread()); + return context; } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandlerTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandlerTest.java index dadffcf..3721589 100644 --- a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandlerTest.java +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/dispatcher/FlowEndNodeHandlerTest.java @@ -52,6 +52,11 @@ public class FlowEndNodeHandlerTest { assertTrue(handler.handle(earlierItem).isSuccess()); assertEquals(0, awaitAllCalls.get()); + TaskNodeExecuteMessage nestedFlowEnd = message(0); + nestedFlowEnd.setNestedFlow(true); + assertTrue(handler.handle(nestedFlowEnd).isSuccess()); + assertEquals(0, awaitAllCalls.get()); + TaskNodeExecuteMessage lastItem = message(0); assertTrue(handler.handle(lastItem).isSuccess()); assertEquals(1, awaitAllCalls.get()); diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutorTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutorTest.java new file mode 100644 index 0000000..dab2c43 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutorTest.java @@ -0,0 +1,66 @@ +package com.cmvr.test.flow.runtime.engine; + +import com.cmvr.test.enums.NodeTypeEnum; +import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +import com.cmvr.test.flow.builder.FlowParamDef; +import com.cmvr.test.flow.context.TaskContext; +import com.cmvr.test.flow.context.TaskContextManager; +import com.cmvr.test.flow.context.TaskInstHolder; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; + +import static org.junit.Assert.assertEquals; + +public class FlowItemExecutorTest { + + @Test + public void nestedPublishedGraphEntersThroughItsStartNode() { + TaskContextManager contextManager = new TaskContextManager(); + TaskInstHolder taskInstHolder = new TaskInstHolder(null, contextManager, null, null); + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + taskInstHolder.register(context); + + FlowGraph graph = new FlowGraph(); + FlowNodeWrapper start = new FlowNodeWrapper(); + start.setNodeId("nested-start"); + start.setNodeName("Start"); + start.setNodeType(NodeTypeEnum.START); + LinkedHashMap nodes = new LinkedHashMap<>(); + nodes.put(start.getNodeId(), start); + graph.setNodeMap(nodes); + graph.setStartInputParamsDefs(new ArrayList<>()); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setInstId(context.getInstId()); + RecordingFlowItemExecutor executor = new RecordingFlowItemExecutor(taskInstHolder); + + executor.executeNestedGraph(graph, message, () -> { }, List.of()); + + assertEquals(start.getNodeId(), executor.executedNodeId); + } + + private static class RecordingFlowItemExecutor extends FlowItemExecutor { + private String executedNodeId; + + private RecordingFlowItemExecutor(TaskInstHolder taskInstHolder) { + super(null, List.of(), null, taskInstHolder, null, null, null, null); + } + + @Override + public void executeNode(FlowGraph graph, + FlowNodeWrapper node, + String startNodeId, + List startInputDefs, + TaskNodeExecuteMessage rootMessage, + Runnable onFinished, + List iterations) { + executedNodeId = node.getNodeId(); + } + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowTaskSchedulerTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowTaskSchedulerTest.java new file mode 100644 index 0000000..b224166 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/FlowTaskSchedulerTest.java @@ -0,0 +1,57 @@ +package com.cmvr.test.flow.runtime.engine; + +import com.cmvr.test.enums.NodeTypeEnum; +import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.builder.FlowNodeWrapper; +import com.cmvr.test.flow.context.TaskContext; +import com.cmvr.test.flow.runtime.engine.support.FlowSchedulerCache; +import org.junit.Test; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +import java.util.LinkedHashMap; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class FlowTaskSchedulerTest { + + @Test + public void subGraphSchedulesEveryIndependentEntryLine() throws Exception { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(2); + executor.setMaxPoolSize(2); + executor.setQueueCapacity(4); + executor.initialize(); + try { + FlowTaskScheduler scheduler = new FlowTaskScheduler(executor, new FlowSchedulerCache()); + FlowGraph graph = new FlowGraph(); + graph.setSub(true); + LinkedHashMap nodes = new LinkedHashMap<>(); + nodes.put("line-a", node("line-a")); + nodes.put("line-b", node("line-b")); + graph.setNodeMap(nodes); + + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + CountDownLatch executed = new CountDownLatch(2); + + scheduler.subStart(graph, context, ignored -> executed.countDown(), () -> { }, List.of(1)); + + assertTrue(executed.await(1, TimeUnit.SECONDS)); + assertEquals(0, executed.getCount()); + } finally { + executor.shutdown(); + } + } + + private FlowNodeWrapper node(String id) { + FlowNodeWrapper node = new FlowNodeWrapper(); + node.setNodeId(id); + node.setNodeName(id); + node.setNodeType(NodeTypeEnum.SLEEP); + return node; + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparerTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparerTest.java new file mode 100644 index 0000000..7e40722 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparerTest.java @@ -0,0 +1,159 @@ +package com.cmvr.test.flow.runtime.engine.support; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.builder.FlowModelBuilder; +import com.cmvr.test.flow.context.TaskContext; +import com.cmvr.test.flow.context.TaskContextManager; +import com.cmvr.test.flow.context.TaskInstHolder; +import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.operator.llm.AsyncMediaAnalysisCoordinator; +import org.junit.Test; + +import java.util.List; + +import static org.junit.Assert.assertEquals; + +public class FlowNodeParamPreparerTest { + + @Test + public void legacyBlankLoopReferenceUsesStartCountParameter() { + FlowGraph graph = graphWithLoopParam("\"quote\":\"\""); + FlowNodeParamPreparer preparer = preparerWithRunParams( + new JSONObject().fluentPut("count", "2")); + + JSONObject input = preparer.prepare(graph, graph.getNode("loop-1"), message()); + + assertEquals(2, input.getIntValue("loopNum")); + assertEquals("2", input.getString("loopArray")); + } + + @Test + public void explicitStartReferenceUsesRuntimeValueInsteadOfDefault() { + FlowGraph graph = graphWithLoopParam( + "\"quote\":[\"start-1\",\"input\",\"count\"],\"quoteType\":\"input\""); + FlowNodeParamPreparer preparer = preparerWithRunParams( + new JSONObject().fluentPut("count", "3")); + + JSONObject input = preparer.prepare(graph, graph.getNode("loop-1"), message()); + + assertEquals(3, input.getIntValue("loopNum")); + assertEquals("3", input.getString("loopArray")); + } + + @Test + public void endNodeCanReturnStartInputAndPreviousNodeOutput() { + String flow = """ + { + "nodes": [ + {"id":"start-1","type":"start","properties":{"name":"Start","action":"START", + "inputParams":[{"name":"requestId","type":"string","input":"default"}]}}, + {"id":"source-1","type":"sleep","properties":{"name":"Source","action":"SLEEP", + "nodeParams":[],"outputParams":[{"name":"value","type":"number"}]}}, + {"id":"end-1","type":"end","properties":{"name":"End","action":"END", + "nodeParams":[ + {"name":"requestId","type":"quote","quote":["start-1","input","requestId"],"quoteType":"input"}, + {"name":"result","type":"quote","quote":["source-1","output","value"],"quoteType":"output"} + ],"outputParams":[{"name":"requestId","type":"string"},{"name":"result","type":"number"}]}} + ], + "edges": [] + } + """; + FlowGraph graph = FlowModelBuilder.buildExecutableGraph(flow); + TaskContextManager manager = new TaskContextManager(); + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + context.setRunParams(new JSONObject().fluentPut("requestId", "runtime-id")); + context.setNodeOutput("source-1", List.of(), new JSONObject().fluentPut("value", 42)); + manager.register(context); + TaskInstHolder holder = new TaskInstHolder( + null, manager, null, new FlowExecutionEventPublisher(List.of())); + AsyncMediaAnalysisCoordinator coordinator = new AsyncMediaAnalysisCoordinator(holder, null, null, null); + FlowNodeParamPreparer preparer = new FlowNodeParamPreparer(holder, coordinator); + + JSONObject input = preparer.prepare(graph, graph.getNode("end-1"), message()); + + assertEquals("runtime-id", input.getString("requestId")); + assertEquals(42, input.getIntValue("result")); + } + + @Test + public void nestedFlowStartUsesSubFlowInputInsteadOfTaskRunParams() { + String flow = """ + { + "nodes": [ + {"id":"start-1","type":"start","properties":{"name":"Start","action":"START", + "inputParams":[{"name":"requestId","type":"string","input":"default"}]}} + ], + "edges": [] + } + """; + FlowGraph graph = FlowModelBuilder.buildExecutableGraph(flow); + TaskContextManager manager = new TaskContextManager(); + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + context.setRunParams(new JSONObject().fluentPut("requestId", "task-value")); + manager.register(context); + TaskInstHolder holder = new TaskInstHolder( + null, manager, null, new FlowExecutionEventPublisher(List.of())); + FlowNodeParamPreparer preparer = new FlowNodeParamPreparer(holder, null); + TaskNodeExecuteMessage message = message(); + message.setLocalRunParams(new JSONObject().fluentPut("requestId", "sub-flow-value")); + + JSONObject input = preparer.prepare(graph, graph.getNode("start-1"), message); + + assertEquals("sub-flow-value", input.getString("requestId")); + } + + private FlowGraph graphWithLoopParam(String quoteFields) { + String flow = """ + { + "nodes": [ + { + "id": "start-1", + "type": "start", + "properties": { + "name": "Start", + "action": "START", + "inputParams": [ + {"name":"count","type":"string","input":"1"} + ] + } + }, + { + "id": "loop-1", + "type": "loop", + "children": [], + "properties": { + "name": "Loop", + "action": "START_LOOP", + "nodeParams": [ + {"name":"loopNum","type":"quote","input":"",%s} + ] + } + } + ], + "edges": [] + } + """.formatted(quoteFields); + return FlowModelBuilder.buildExecutableGraph(flow); + } + + private FlowNodeParamPreparer preparerWithRunParams(JSONObject runParams) { + TaskContextManager manager = new TaskContextManager(); + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + context.setRunParams(runParams); + manager.register(context); + TaskInstHolder holder = new TaskInstHolder( + null, manager, null, new FlowExecutionEventPublisher(List.of())); + return new FlowNodeParamPreparer(holder, null); + } + + private TaskNodeExecuteMessage message() { + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setInstId("inst-1"); + return message; + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java index dce68a3..42018a7 100644 --- a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java @@ -2,6 +2,8 @@ package com.cmvr.test.flow.runtime.engine.support; import org.junit.Test; +import java.util.Set; + import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -30,4 +32,14 @@ public class FlowSchedulerCacheTest { assertFalse(cache.trySchedule(otherInstanceKey)); } + + @Test + public void subGraphCompletesOnlyAfterAllTerminalNodes() { + FlowSchedulerCache cache = new FlowSchedulerCache(); + String key = "inst-1:subgraph-terminal:loop-1:1"; + + assertFalse(cache.markGraphTerminalCompleted(key, "audio-end", Set.of("audio-end", "video-end"))); + assertTrue(cache.markGraphTerminalCompleted(key, "video-end", Set.of("audio-end", "video-end"))); + assertFalse(cache.markGraphTerminalCompleted(key, "video-end", Set.of("audio-end", "video-end"))); + } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptorTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptorTest.java new file mode 100644 index 0000000..12efe32 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/interceptor/FlowBeforeInterceptorTest.java @@ -0,0 +1,59 @@ +package com.cmvr.test.flow.runtime.interceptor; + +import com.cmvr.test.flow.builder.FlowGraph; +import com.cmvr.test.flow.context.TaskContext; +import com.cmvr.test.flow.context.TaskContextManager; +import com.cmvr.test.flow.context.TaskInstHolder; +import com.cmvr.test.flow.runtime.event.FlowExecutionEvent; +import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher; +import com.cmvr.test.flow.runtime.event.FlowExecutionListener; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.Assert.assertEquals; + +public class FlowBeforeInterceptorTest { + + @Test + public void publishesLoopIterationsOnNodeStarted() { + AtomicReference published = new AtomicReference<>(); + FlowExecutionListener listener = new FlowExecutionListener() { + @Override + public void onEvent(FlowExecutionEvent event) { + published.set(event); + } + + @Override + public String getModuleCode() { + return "AIMA"; + } + }; + FlowExecutionEventPublisher publisher = new FlowExecutionEventPublisher(List.of(listener)); + TaskContextManager contextManager = new TaskContextManager(); + TaskContext context = new TaskContext(); + context.setInstId("inst-1"); + context.setModuleCode("AIMA"); + contextManager.register(context); + TaskInstHolder holder = new TaskInstHolder(null, contextManager, null, publisher); + FlowBeforeInterceptor interceptor = new FlowBeforeInterceptor(publisher, holder); + + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setInstId("inst-1"); + message.setNodeId("node-1"); + message.setNodeName("循环内节点"); + message.setNodeType("FUNCTION"); + message.setGraph(new FlowGraph()); + message.setIterations(new ArrayList<>(List.of(2, 3))); + + interceptor.doIntercept(message, ignored -> TaskNodeExecuteResult.success()); + message.getIterations().clear(); + + assertEquals(FlowExecutionEvent.EventType.NODE_STARTED, published.get().getEventType()); + assertEquals(List.of(2, 3), published.get().getIterations()); + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinatorTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinatorTest.java index 44bebbb..2274379 100644 --- a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinatorTest.java +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/AsyncMediaAnalysisCoordinatorTest.java @@ -17,6 +17,7 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.lang.reflect.Proxy; import java.util.Collections; import java.util.List; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -79,6 +80,71 @@ public class AsyncMediaAnalysisCoordinatorTest { } } + @Test + public void reportsBackgroundFailureAtOwningSubFlowBarrier() { + Fixture fixture = fixture(); + try { + fixture.message.setSubFlowExecutionId("sub-flow-scope"); + fixture.coordinator.submit(fixture.message, "request-sub-flow", () -> { + throw new IllegalStateException("nested vision failed"); + }); + + TaskNodeExecuteResult result = fixture.coordinator.awaitScope( + "inst-1", "sub-flow-scope"); + + assertFalse(result.isSuccess()); + assertTrue(result.getErrorMsg().contains("nested vision failed")); + } finally { + fixture.executor.shutdown(); + } + } + + @Test + public void cancelledSubFlowScopeCannotPublishLateResult() throws Exception { + Fixture fixture = fixture(); + CountDownLatch started = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + try { + fixture.message.setSubFlowExecutionId("cancelled-scope"); + fixture.coordinator.submit(fixture.message, "request-cancelled", () -> { + started.countDown(); + await(release); + return TaskNodeExecuteResult.success(new JSONObject().fluentPut("passed", true)); + }); + assertTrue(started.await(1, TimeUnit.SECONDS)); + + fixture.coordinator.cancelScope("inst-1", "cancelled-scope"); + release.countDown(); + Thread.sleep(50); + + assertNull(fixture.context.getNodeOutput("analysis-node", Collections.emptyList())); + } finally { + release.countDown(); + fixture.executor.shutdown(); + } + } + + @Test + public void cancellingScopeRemovesFailureFromFinalTaskBarrier() throws Exception { + Fixture fixture = fixture(); + CountDownLatch failed = new CountDownLatch(1); + try { + fixture.message.setSubFlowExecutionId("failed-old-scope"); + fixture.coordinator.submit(fixture.message, "request-old-failure", () -> { + failed.countDown(); + throw new IllegalStateException("failure before pause"); + }); + assertTrue(failed.await(1, TimeUnit.SECONDS)); + fixture.coordinator.awaitScope("inst-1", "failed-old-scope"); + + fixture.coordinator.cancelScope("inst-1", "failed-old-scope"); + + assertTrue(fixture.coordinator.awaitAll("inst-1").isSuccess()); + } finally { + fixture.executor.shutdown(); + } + } + @Test public void unrelatedItemFailureDoesNotBreakCurrentItemDependencyWait() { Fixture fixture = fixture(); @@ -100,6 +166,39 @@ public class AsyncMediaAnalysisCoordinatorTest { } } + @Test + public void repeatedNodeKeyDoesNotReplaceAnUnfinishedAnalysis() throws Exception { + Fixture fixture = fixture(); + CountDownLatch firstStarted = new CountDownLatch(1); + CountDownLatch releaseFirst = new CountDownLatch(1); + CountDownLatch secondFinished = new CountDownLatch(1); + try { + fixture.coordinator.submit(fixture.message, "request-first", () -> { + firstStarted.countDown(); + await(releaseFirst); + return TaskNodeExecuteResult.success(); + }); + assertTrue(firstStarted.await(1, TimeUnit.SECONDS)); + + fixture.coordinator.submit(fixture.message, "request-second", () -> { + secondFinished.countDown(); + return TaskNodeExecuteResult.success(); + }); + assertTrue(secondFinished.await(1, TimeUnit.SECONDS)); + + CompletableFuture barrier = CompletableFuture.supplyAsync( + () -> fixture.coordinator.awaitAll("inst-1")); + Thread.sleep(100); + assertFalse(barrier.isDone()); + + releaseFirst.countDown(); + assertTrue(barrier.get(1, TimeUnit.SECONDS).isSuccess()); + } finally { + releaseFirst.countDown(); + fixture.executor.shutdown(); + } + } + private Fixture fixture() { AtomicBoolean nodeSucceeded = new AtomicBoolean(); AtomicBoolean nodeFailed = new AtomicBoolean(); @@ -130,8 +229,8 @@ public class AsyncMediaAnalysisCoordinatorTest { contextManager.register(context); TaskInstHolder holder = new TaskInstHolder(nodeService, contextManager, null, publisher); ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); - executor.setCorePoolSize(1); - executor.setMaxPoolSize(1); + executor.setCorePoolSize(2); + executor.setMaxPoolSize(2); executor.setQueueCapacity(10); executor.setThreadNamePrefix("media-analysis-test-"); executor.initialize(); diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionServiceTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionServiceTest.java new file mode 100644 index 0000000..cb6b0e6 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/FlowSubFlowDefinitionServiceTest.java @@ -0,0 +1,253 @@ +package com.cmvr.test.flow.runtime.subflow; + +import com.cmvr.test.mapper.TeDetectionItemMapper; +import com.cmvr.test.mapper.TeDetectionItemVersionMapper; +import com.cmvr.test.model.domain.TeDetectionItem; +import com.cmvr.test.model.domain.TeDetectionItemVersion; +import com.cmvr.common.exception.GlobalException; +import org.junit.Test; + +import java.lang.reflect.Proxy; +import java.util.List; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThrows; +import static org.junit.Assert.assertTrue; + +public class FlowSubFlowDefinitionServiceTest { + + @Test + public void loadsImmutableVersionAndNamespacesInternalReferences() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem child = item("child", "Child", "child-v1"); + TeDetectionItemVersion version = new TeDetectionItemVersion(); + version.setId("child-v1"); + version.setDetectionItemId("child"); + version.setVersionNo(1); + version.setFlowData(childFlow()); + + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapper(parent, child), versionMapper(version)); + + SubFlowExecutionDefinition definition = service.loadExecution( + "parent", "child", "child-v1", "parent-node"); + + assertTrue(definition.endNodeId().startsWith("sf_parent-node_")); + assertTrue(definition.endNodeId().length() <= 32); + assertEquals(1, definition.graph().getInDegreeZeroNodes().size()); + assertTrue(definition.graph().allNodeIds().stream().allMatch(id -> id.length() <= 32)); + assertEquals(definition.graph().allNodeIds().size(), + definition.graph().allNodeIds().stream().distinct().count()); + String startNodeId = definition.graph().getNodeMap().values().stream() + .filter(node -> node.getNodeType() == com.cmvr.test.enums.NodeTypeEnum.START) + .map(node -> node.getNodeId()) + .findFirst() + .orElseThrow(); + assertNotNull(definition.graph().getNode(startNodeId)); + assertNotNull(definition.graph().getNode(definition.endNodeId())); + assertEquals(startNodeId, + definition.graph().getNode(definition.endNodeId()) + .getNodeParams().get(0).getQuote().get(0)); + } + + @Test + public void optionsOnlyContainPublishedWorkflowsFromSameProject() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem sameProject = item("child", "Child", "child-v1"); + TeDetectionItem aima = item("aima", "Aima", "aima-v1"); + TeDetectionItemVersion childVersion = version("child-v1", "child"); + TeDetectionItemVersion aimaVersion = version("aima-v1", "aima"); + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapperWithList(parent, sameProject, aima), versionMapper(childVersion, aimaVersion)); + + var options = service.listOptions("parent"); + + assertEquals(1, options.size()); + assertEquals("child", options.get(0).getItemId()); + assertEquals("CHANGAN", options.get(0).getProjectCode()); + } + + @Test + public void optionsExcludeWorkflowThatWouldReferenceCurrentWorkflowTransitively() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem child = item("child", "Child", "child-v1"); + TeDetectionItemVersion childVersion = version("child-v1", "child"); + childVersion.setFlowData(flowCalling("parent", "parent-v1")); + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapperWithList(parent, child, item("aima", "Aima", "")), + versionMapper(childVersion)); + + var options = service.listOptions("parent"); + + assertTrue(options.isEmpty()); + } + + @Test + public void publishRejectsTransitiveSubFlowCircularReference() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem child = item("child", "Child", "child-v1"); + TeDetectionItemVersion childVersion = version("child-v1", "child"); + childVersion.setFlowData(flowCalling("parent", "parent-v1")); + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapper(parent, child), versionMapper(childVersion)); + + GlobalException error = assertThrows(GlobalException.class, + () -> service.validateForPublish(parent, flowCalling("child", "child-v1"))); + + assertEquals("检测到子流程循环引用:Parent -> Child -> Parent", error.getMessage()); + } + + @Test + public void publishRejectsSubFlowInputThatReferencesDownstreamNode() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem child = item("child", "Child", "child-v1"); + TeDetectionItemVersion version = version("child-v1", "child"); + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapper(parent, child), versionMapper(version)); + + GlobalException error = assertThrows(GlobalException.class, + () -> service.validateForPublish(parent, parentFlow("later"))); + + assertEquals("子流程节点[SubFlow]的入参[value]引用的节点不是当前节点的上游节点", error.getMessage()); + } + + @Test + public void publishAcceptsEndOutputReferencingUpstreamSubFlowOutput() { + TeDetectionItem parent = item("parent", "Parent", "parent-v1"); + TeDetectionItem child = item("child", "Child", "child-v1"); + TeDetectionItemVersion version = version("child-v1", "child"); + FlowSubFlowDefinitionService service = new FlowSubFlowDefinitionService( + itemMapper(parent, child), versionMapper(version)); + + service.validateForPublish(parent, parentFlow("start")); + } + + private TeDetectionItem item(String id, String name, String versionId) { + TeDetectionItem item = new TeDetectionItem(); + item.setId(id); + item.setDetectName(name); + item.setPublishedVersionId(versionId); + item.setStatus("0"); + return item; + } + + private TeDetectionItemVersion version(String id, String itemId) { + TeDetectionItemVersion version = new TeDetectionItemVersion(); + version.setId(id); + version.setDetectionItemId(itemId); + version.setVersionNo(1); + version.setFlowData(childFlow()); + return version; + } + + private TeDetectionItemMapper itemMapper(TeDetectionItem... items) { + return mapper(TeDetectionItemMapper.class, (method, args) -> { + if ("selectAimaDetectionItemIds".equals(method)) return List.of(); + if ("selectById".equals(method)) { + for (TeDetectionItem item : items) if (item.getId().equals(args[0])) return item; + } + return defaultValue(method); + }); + } + + private TeDetectionItemMapper itemMapperWithList(TeDetectionItem parent, TeDetectionItem child, TeDetectionItem aima) { + return mapper(TeDetectionItemMapper.class, (method, args) -> { + if ("selectAimaDetectionItemIds".equals(method)) return List.of("aima"); + if ("selectById".equals(method)) { + if (parent.getId().equals(args[0])) return parent; + if (child.getId().equals(args[0])) return child; + if (aima.getId().equals(args[0])) return aima; + } + if ("selectList".equals(method)) return List.of(child, aima); + return defaultValue(method); + }); + } + + private TeDetectionItemVersionMapper versionMapper(TeDetectionItemVersion... versions) { + return mapper(TeDetectionItemVersionMapper.class, (method, args) -> { + if ("selectById".equals(method)) { + for (TeDetectionItemVersion version : versions) if (version.getId().equals(args[0])) return version; + } + return defaultValue(method); + }); + } + + @SuppressWarnings("unchecked") + private T mapper(Class type, Invocation invocation) { + return (T) Proxy.newProxyInstance(type.getClassLoader(), new Class[]{type}, + (proxy, method, args) -> invocation.call(method.getName(), args)); + } + + private Object defaultValue(String method) { + if (method.startsWith("select")) return null; + return 0; + } + + private String childFlow() { + return """ + { + "nodes": [ + {"id":"start","type":"start","properties":{"name":"Start","action":"START", + "inputParams":[{"name":"value","type":"string","input":""}],"resourceRoles":[]}}, + {"id":"end","type":"end","properties":{"name":"End","action":"END", + "nodeParams":[{"name":"result","type":"quote","quote":["start","input","value"],"quoteType":"input"}], + "outputParams":[{"name":"result","type":"string"}]}} + ], + "edges":[{"id":"edge","sourceNodeId":"start","targetNodeId":"end","sourceAnchorId":"start_1","targetAnchorId":"end_3"}] + } + """; + } + + private String parentFlow(String subFlowInputNodeId) { + return """ + { + "nodes": [ + {"id":"start","type":"start","properties":{"name":"Start","action":"START", + "inputParams":[{"name":"value","type":"string","input":""}],"resourceRoles":[]}}, + {"id":"sub","type":"subFlow","properties":{"name":"SubFlow","action":"SUB_FLOW", + "subFlow":{"itemId":"child","versionId":"child-v1","versionNo":1}, + "nodeParams":[{"name":"value","type":"quote","quote":["%s","%s","value"],"quoteType":"%s"}], + "outputParams":[{"name":"result","type":"string"}]}}, + {"id":"later","type":"code","properties":{"name":"Later","action":"EXECUTE_CODE", + "outputParams":[{"name":"value","type":"string"}]}}, + {"id":"end","type":"end","properties":{"name":"End","action":"END", + "nodeParams":[{"name":"result","type":"quote","quote":["sub","output","result"],"quoteType":"output"}], + "outputParams":[{"name":"result","type":"string"}]}} + ], + "edges":[ + {"id":"e1","sourceNodeId":"start","targetNodeId":"sub"}, + {"id":"e2","sourceNodeId":"sub","targetNodeId":"later"}, + {"id":"e3","sourceNodeId":"later","targetNodeId":"end"} + ] + } + """.formatted(subFlowInputNodeId, + "start".equals(subFlowInputNodeId) ? "input" : "output", + "start".equals(subFlowInputNodeId) ? "input" : "output"); + } + + private String flowCalling(String childItemId, String childVersionId) { + return """ + { + "nodes": [ + {"id":"start","type":"start","properties":{"name":"Start","action":"START", + "inputParams":[],"resourceRoles":[]}}, + {"id":"sub","type":"subFlow","properties":{"name":"SubFlow","action":"SUB_FLOW", + "subFlow":{"itemId":"%s","versionId":"%s","versionNo":1}, + "nodeParams":[],"outputParams":[]}}, + {"id":"end","type":"end","properties":{"name":"End","action":"END", + "nodeParams":[],"outputParams":[]}} + ], + "edges":[ + {"id":"e1","sourceNodeId":"start","targetNodeId":"sub"}, + {"id":"e2","sourceNodeId":"sub","targetNodeId":"end"} + ] + } + """.formatted(childItemId, childVersionId); + } + + @FunctionalInterface + private interface Invocation { + Object call(String method, Object[] args); + } +} diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcomeTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcomeTest.java new file mode 100644 index 0000000..c2096b8 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/subflow/SubFlowExecutionOutcomeTest.java @@ -0,0 +1,20 @@ +package com.cmvr.test.flow.runtime.subflow; + +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class SubFlowExecutionOutcomeTest { + + @Test + public void retainsFirstNestedFailure() { + SubFlowExecutionOutcome outcome = new SubFlowExecutionOutcome(); + + outcome.fail("机械臂移动", "执行异常:设备离线"); + outcome.fail("结束节点", "不应覆盖首个错误"); + + assertTrue(outcome.isFailed()); + assertEquals("内部节点[机械臂移动]执行异常:设备离线", outcome.getFailure()); + } +}