feat(flow): 添加子流程功能支持和循环迭代追踪
- 新增 SUB_FLOW 类型枚举值用于标识子流程节点 - 实现子流程异步媒体分析协调器管理机制 - 添加循环迭代路径追踪到执行上下文日志 - 增强流程控制服务支持嵌套子流程暂停恢复 - 实现子流程作用域内异步任务等待和取消功能 - 优化循环体多入口独立执行线程调度机制 - 添加节点参数循环引用检测验证机制 - 支持遗留循环参数值解析兼容性处理 - 配置虚拟机器人ID列表绕过在线检查机制
This commit is contained in:
parent
6e01f6888b
commit
0948d73a62
@ -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) {
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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)
|
||||
|
||||
@ -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<Integer> 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();
|
||||
}
|
||||
|
||||
@ -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));
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<String> virtualRobotIds = new ArrayList<>();
|
||||
|
||||
public long getHeartbeatTimeoutMs() {
|
||||
return heartbeatTimeoutMs;
|
||||
}
|
||||
|
||||
public void setHeartbeatTimeoutMs(long heartbeatTimeoutMs) {
|
||||
this.heartbeatTimeoutMs = heartbeatTimeoutMs;
|
||||
}
|
||||
|
||||
public List<String> getVirtualRobotIds() {
|
||||
return virtualRobotIds;
|
||||
}
|
||||
|
||||
public void setVirtualRobotIds(List<String> 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<String> 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);
|
||||
}
|
||||
}
|
||||
@ -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");
|
||||
|
||||
@ -35,7 +35,15 @@ public interface InspectionRobotMapper extends MPJBaseMapper<InspectionRobot>
|
||||
/**
|
||||
* Select connected robots that currently report one online AGV device.
|
||||
*/
|
||||
List<InspectionRobot> selectRuntimeSyncCandidates();
|
||||
List<InspectionRobot> selectRuntimeSyncCandidates(@Param("heartbeatCutoff") Date heartbeatCutoff,
|
||||
@Param("excludedRobotIds") List<String> excludedRobotIds);
|
||||
|
||||
List<InspectionRobot> selectStaleConnectedRobots(@Param("heartbeatCutoff") Date heartbeatCutoff,
|
||||
@Param("excludedRobotIds") List<String> 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);
|
||||
|
||||
|
||||
@ -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 + "]");
|
||||
|
||||
@ -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<InspectionRobot> 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);
|
||||
}
|
||||
|
||||
@ -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);
|
||||
if (!robotStateProperties.isVirtualRobot(normalized)) {
|
||||
endpointResolver.resolve(normalized);
|
||||
}
|
||||
return new RobotExecutionTarget(normalized);
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<InspectionRobot> 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());
|
||||
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<String> 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<String> reportedIds = new ArrayList<>();
|
||||
for (ManagedDeviceStatus status : node.getDeviceManager().getDevicesList()) {
|
||||
|
||||
@ -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();
|
||||
}
|
||||
}
|
||||
@ -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<String> 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;
|
||||
}
|
||||
|
||||
@ -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}
|
||||
<if test="excludedRobotIds != null and excludedRobotIds.size() > 0">
|
||||
and lower(r.robot_id) not in
|
||||
<foreach collection="excludedRobotIds" item="robotId" open="(" separator="," close=")">
|
||||
#{robotId}
|
||||
</foreach>
|
||||
</if>
|
||||
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
|
||||
</select>
|
||||
|
||||
<select id="selectStaleConnectedRobots" resultMap="InspectionRobotResult">
|
||||
select r.id, r.robot_id, r.connect_status, r.last_heartbeat_time, r.archived_status
|
||||
from inspection_robot r
|
||||
where r.connect_status = '1'
|
||||
and coalesce(r.archived_status, '0') = '0'
|
||||
and (r.last_heartbeat_time is null or r.last_heartbeat_time < #{heartbeatCutoff})
|
||||
<if test="excludedRobotIds != null and excludedRobotIds.size() > 0">
|
||||
and lower(r.robot_id) not in
|
||||
<foreach collection="excludedRobotIds" item="robotId" open="(" separator="," close=")">
|
||||
#{robotId}
|
||||
</foreach>
|
||||
</if>
|
||||
</select>
|
||||
|
||||
<update id="markHeartbeatExpiredOffline">
|
||||
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>
|
||||
|
||||
<update id="archiveInspectionRobotByIds" parameterType="String">
|
||||
update inspection_robot set archived_status = '1', connect_status = '0',
|
||||
update_time = sysdate(), update_by = #{updateBy}
|
||||
|
||||
@ -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<NodeEvent> listener) {
|
||||
if (closed.get()) {
|
||||
|
||||
@ -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();
|
||||
|
||||
@ -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", "子流程结束"),
|
||||
|
||||
@ -15,6 +15,7 @@ public enum NodeTypeEnum {
|
||||
FUNCTION("function"),
|
||||
BRANCH("branch"),
|
||||
LOOP("loop"),
|
||||
SUB_FLOW("subFlow"),
|
||||
STOP_LOOP("stopLoop"),
|
||||
END("end"),
|
||||
SLEEP("sleep"),
|
||||
|
||||
@ -114,15 +114,25 @@ public class FlowGraph {
|
||||
* 获取子图中起始节点(入度为0)
|
||||
*/
|
||||
public FlowNodeWrapper getInDegreeZeroNode() {
|
||||
List<FlowNodeWrapper> 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<FlowNodeWrapper> getInDegreeZeroNodes() {
|
||||
lock.readLock().lock();
|
||||
try {
|
||||
List<FlowNodeWrapper> result = new ArrayList<>();
|
||||
for (FlowNodeWrapper node : nodeMap.values()) {
|
||||
List<FlowEdge> flowEdges = predecessors.get(node.getNodeId());
|
||||
if (flowEdges == null || flowEdges.isEmpty()) {
|
||||
return node;
|
||||
result.add(node);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
return result;
|
||||
} finally {
|
||||
lock.readLock().unlock();
|
||||
}
|
||||
|
||||
@ -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<String, Set<String>> dependencies = new LinkedHashMap<>();
|
||||
Map<String, String> 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<String> references = new HashSet<>();
|
||||
collectActiveReferences(properties, references);
|
||||
dependencies.put(nodeId, references);
|
||||
nodeNames.put(nodeId, properties == null
|
||||
? nodeId : StrUtil.blankToDefault(properties.getString("name"), nodeId));
|
||||
}
|
||||
|
||||
Map<String, Integer> states = new HashMap<>();
|
||||
List<String> path = new ArrayList<>();
|
||||
for (String nodeId : dependencies.keySet()) {
|
||||
List<String> cycle = findReferenceCycle(nodeId, dependencies, states, path);
|
||||
if (cycle == null) continue;
|
||||
List<String> names = cycle.stream().map(id -> nodeNames.getOrDefault(id, id)).toList();
|
||||
throw new GlobalException("节点参数存在循环引用:" + String.join(" -> ", names));
|
||||
}
|
||||
}
|
||||
|
||||
private static void collectActiveReferences(Object value, Set<String> 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<String> references) {
|
||||
if (quote != null && quote.size() >= 3 && StrUtil.isNotBlank(quote.getString(0))) {
|
||||
references.add(quote.getString(0));
|
||||
}
|
||||
}
|
||||
|
||||
private static List<String> findReferenceCycle(String nodeId,
|
||||
Map<String, Set<String>> dependencies,
|
||||
Map<String, Integer> states,
|
||||
List<String> 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<String> 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<String> 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<String, FlowNodeWrapper> nodeMap, FlowGraph graph) {
|
||||
Map<String, List<FlowParamDef>> nodeParams = new LinkedHashMap<>();
|
||||
for (Map.Entry<String, FlowNodeWrapper> 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":
|
||||
|
||||
@ -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<Integer> 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) {
|
||||
|
||||
@ -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<TaskNodeExecuteContext> discardNestedSubFlowContexts(String instId) {
|
||||
List<TaskNodeExecuteContext> removed = new ArrayList<>();
|
||||
for (Map.Entry<String, TaskNodeExecuteContext> 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)
|
||||
*/
|
||||
|
||||
@ -220,8 +220,22 @@ public class FlowControlService {
|
||||
}
|
||||
}
|
||||
|
||||
List<TaskNodeExecuteContext> 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<TaskNodeExecuteContext> discardedSubFlowNodes =
|
||||
taskThreadRegistry.discardNestedSubFlowContexts(instId);
|
||||
for (TaskNodeExecuteContext discarded : discardedSubFlowNodes) {
|
||||
nodeInstService.delete(instId, discarded.getNode().getNodeId(), discarded.getIterations());
|
||||
}
|
||||
ctx.setPaused(false);
|
||||
|
||||
List<TaskNodeExecuteContext> pendingNodes = taskThreadRegistry.getPendingNodes(instId);
|
||||
List<TaskNodeExecuteContext> pendingLoops = pendingNodes.stream()
|
||||
.filter(node -> node.getNode().getNodeType().equals(NodeTypeEnum.LOOP))
|
||||
|
||||
@ -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;
|
||||
|
||||
@ -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<String> 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);
|
||||
}
|
||||
}
|
||||
@ -81,18 +81,19 @@ public class FlowItemExecutor {
|
||||
// 主图中开始节点的 inputParams
|
||||
List<FlowParamDef> 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<FlowNodeWrapper> 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<Integer> iterations) {
|
||||
String startNodeId = graph.findStartNodeId();
|
||||
List<FlowParamDef> 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,6 +169,9 @@ public class FlowItemExecutor {
|
||||
|
||||
// END 节点 or 子图结束
|
||||
if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END || graph.isSubGraphEndNode(nodeId)) {
|
||||
boolean graphFinished = !graph.isSub()
|
||||
|| flowTaskScheduler.markSubGraphTerminalCompleted(instId, nodeId, graph, iterations);
|
||||
if (graphFinished) {
|
||||
try {
|
||||
log.info("{} 节点执行,释放线程", node.getNodeType());
|
||||
onFinished.run();
|
||||
@ -155,6 +179,7 @@ public class FlowItemExecutor {
|
||||
log.error("onFinished 回调异常", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 通知调度器继续调度后继节点
|
||||
flowTaskScheduler.markCompleted(
|
||||
@ -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()) {
|
||||
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, "节点执行异常: " + e.getMessage(), iterations);
|
||||
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);
|
||||
|
||||
@ -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<Integer> iterations) {
|
||||
String instId = context.getInstId();
|
||||
FlowNodeWrapper startNode = graph.getInDegreeZeroNode();
|
||||
List<String> nextNodes = graph.getNextNodes(startNode.getNodeId());
|
||||
List<FlowNodeWrapper> 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<Integer> iterations) {
|
||||
Set<String> 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;
|
||||
}
|
||||
/**
|
||||
* 提交节点执行任务(除循环只调度一次)
|
||||
*/
|
||||
|
||||
@ -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<FlowParamDef> 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<FlowParamDef> 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<FlowParamDef> 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
|
||||
|
||||
@ -13,6 +13,8 @@ public class FlowSchedulerCache {
|
||||
// instId_nodeId 或 instId_nodeId_loopIteration
|
||||
private final ConcurrentMap<String, AtomicBoolean> scheduled = new ConcurrentHashMap<>();
|
||||
private final ConcurrentMap<String, Set<String>> completedPreMap = new ConcurrentHashMap<>();
|
||||
private final ConcurrentMap<String, Set<String>> completedTerminalMap = new ConcurrentHashMap<>();
|
||||
private final ConcurrentMap<String, AtomicBoolean> 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<String> terminalNodeIds) {
|
||||
Set<String> 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));
|
||||
}
|
||||
}
|
||||
|
||||
@ -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);
|
||||
|
||||
@ -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);
|
||||
|
||||
@ -28,6 +28,7 @@ public class TaskNodeExecuteContext {
|
||||
private int loopIteration;
|
||||
private Object loopArray;
|
||||
private List<Integer> iterations = new ArrayList<>();
|
||||
private String activeSubFlowExecutionId;
|
||||
|
||||
public void markExecutionFinished() {
|
||||
executionFinished.countDown();
|
||||
|
||||
@ -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<String> 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;
|
||||
|
||||
/**
|
||||
* 上游输出参数(来自上一个节点的输出)
|
||||
*/
|
||||
|
||||
@ -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<Void> job = group.jobs.get(key);
|
||||
CompletableFuture<Void> 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<String, CompletableFuture<Void>> 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<String, CompletableFuture<Void>> scopedJobs = group.scopeJobs.remove(executionId);
|
||||
group.scopeFailures.remove(executionId);
|
||||
Set<String> 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<TaskNodeExecuteResult> analysis) {
|
||||
try {
|
||||
analysis.get();
|
||||
@ -220,7 +285,13 @@ public class AsyncMediaAnalysisCoordinator {
|
||||
|
||||
private static final class ExecutionGroup {
|
||||
private final Map<String, CompletableFuture<Void>> jobs = new ConcurrentHashMap<>();
|
||||
private final Map<String, CompletableFuture<Void>> nodeJobs = new ConcurrentHashMap<>();
|
||||
private final Map<String, String> nodeFailures = new ConcurrentHashMap<>();
|
||||
private final Map<String, Map<String, CompletableFuture<Void>>> scopeJobs = new ConcurrentHashMap<>();
|
||||
private final Map<String, String> scopeFailures = new ConcurrentHashMap<>();
|
||||
private final Map<String, Set<String>> scopeNodeKeys = new ConcurrentHashMap<>();
|
||||
private final Set<String> cancelledScopes = ConcurrentHashMap.newKeySet();
|
||||
private final AtomicLong sequence = new AtomicLong();
|
||||
private final AtomicReference<String> failure = new AtomicReference<>();
|
||||
private volatile boolean cancelled;
|
||||
}
|
||||
|
||||
@ -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<TeSubFlowOptionVO> listOptions(String currentItemId) {
|
||||
TeDetectionItem current = requireItem(currentItemId);
|
||||
Set<String> aimaItemIds = new HashSet<>(itemMapper.selectAimaDetectionItemIds());
|
||||
String projectCode = projectCode(current.getId(), aimaItemIds);
|
||||
List<TeDetectionItem> items = itemMapper.selectList(new LambdaQueryWrapper<TeDetectionItem>()
|
||||
.isNotNull(TeDetectionItem::getPublishedVersionId)
|
||||
.ne(TeDetectionItem::getPublishedVersionId, "")
|
||||
.ne(TeDetectionItem::getId, currentItemId)
|
||||
.eq(TeDetectionItem::getStatus, "0"));
|
||||
List<TeSubFlowOptionVO> 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<String> 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<String> 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<String> 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<String> cycle = new ArrayList<>(path.subList(repeatedAt, path.size()));
|
||||
cycle.add(childItemId);
|
||||
return new CircularReference(cycle, false);
|
||||
}
|
||||
if (depth >= MAX_SUB_FLOW_DEPTH) {
|
||||
List<String> 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<String> itemIds) {
|
||||
List<String> 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<String> 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<String> 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<String, JSONObject> 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<String, JSONObject> paramsByName = byName(targetProperties == null
|
||||
? null : targetProperties.getJSONArray("nodeParams"));
|
||||
Set<String> upstreamIds = upstreamNodeIds(flow, targetNode.getString("id"));
|
||||
Map<String, JSONObject> 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<String> upstreamNodeIds(JSONObject flow, String targetNodeId) {
|
||||
Set<String> 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<String> 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<String, JSONObject> nodesById(JSONObject flow) {
|
||||
Map<String, JSONObject> 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<String, JSONObject> 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<String, JSONObject> 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<String, String> 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<String, String> 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<String, String> 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<String, JSONObject> byName(JSONArray values) {
|
||||
return byKey(values, "name");
|
||||
}
|
||||
|
||||
private Map<String, JSONObject> byKey(JSONArray values, String key) {
|
||||
Map<String, JSONObject> 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<String> contractSignature(JSONArray values) {
|
||||
List<String> 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<String> aimaItemIds) {
|
||||
return aimaItemIds.contains(itemId) ? PROJECT_AIMA : PROJECT_CHANGAN;
|
||||
}
|
||||
}
|
||||
@ -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) {
|
||||
}
|
||||
@ -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<String> 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();
|
||||
}
|
||||
}
|
||||
@ -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<TeDetectionItem> {
|
||||
|
||||
@Select("select distinct detect_item_id from aima_test_case where detect_item_id is not null")
|
||||
List<String> selectAimaDetectionItemIds();
|
||||
}
|
||||
|
||||
@ -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;
|
||||
}
|
||||
@ -92,6 +92,7 @@ public class ITeNodeInstServiceImpl extends ServiceImpl<TeNodeInstMapper, TeNode
|
||||
update.eq(TeNodeInst::getInstId, instId)
|
||||
.eq(TeNodeInst::getItemId, itemId)
|
||||
.eq(TeNodeInst::getNodeId, nodeId)
|
||||
.eq(CollUtil.isNotEmpty(iterations), TeNodeInst::getIteration, getIteration(iterations))
|
||||
.set(TeNodeInst::getParamsOut, params)
|
||||
.set(TeNodeInst::getStatus, TaskStatusEnum.FAILED.name())
|
||||
.set(TeNodeInst::getMessage, message)
|
||||
@ -113,6 +114,7 @@ public class ITeNodeInstServiceImpl extends ServiceImpl<TeNodeInstMapper, TeNode
|
||||
inst.setStatus(TaskStatusEnum.FAILED.name());
|
||||
inst.setStartTime(now); // 插入时设置 start 和 end
|
||||
inst.setEndTime(now);
|
||||
inst.setIteration(getIteration(iterations));
|
||||
|
||||
this.save(inst);
|
||||
}
|
||||
|
||||
@ -10,6 +10,7 @@ import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.test.mapper.TeDetectionItemMapper;
|
||||
import com.cmvr.test.mapper.TeDetectionItemVersionMapper;
|
||||
import com.cmvr.test.flow.runtime.subflow.FlowSubFlowDefinitionService;
|
||||
import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.domain.TeDetectionItemVersion;
|
||||
import com.cmvr.test.model.vo.TeFlowPublishVO;
|
||||
@ -27,6 +28,7 @@ public class TeDetectionItemVersionServiceImpl
|
||||
implements ITeDetectionItemVersionService {
|
||||
|
||||
private final TeDetectionItemMapper detectionItemMapper;
|
||||
private final FlowSubFlowDefinitionService subFlowDefinitionService;
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
@ -36,6 +38,7 @@ public class TeDetectionItemVersionServiceImpl
|
||||
if (StrUtil.isBlank(item.getFlowData())) {
|
||||
throw new GlobalException("请先保存工作流草稿");
|
||||
}
|
||||
subFlowDefinitionService.validateForPublish(item, item.getFlowData());
|
||||
String hash = contentHash(item.getFlowData(), item.getConfig());
|
||||
TeDetectionItemVersion current = StrUtil.isBlank(item.getPublishedVersionId())
|
||||
? null : baseMapper.selectById(item.getPublishedVersionId());
|
||||
|
||||
@ -1,9 +1,11 @@
|
||||
package com.cmvr.test.flow.builder;
|
||||
|
||||
import com.cmvr.test.enums.ActionEnum;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertThrows;
|
||||
|
||||
public class FlowModelBuilderTest {
|
||||
|
||||
@ -31,4 +33,24 @@ public class FlowModelBuilderTest {
|
||||
|
||||
assertEquals(ActionEnum.IMAGE_ANALYZE, graph.getNodeMap().get("image-1").getAction());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void rejectsCircularParameterReferencesBeforeExecution() {
|
||||
String flow = """
|
||||
{
|
||||
"nodes": [
|
||||
{"id":"a","type":"code","properties":{"name":"A","action":"EXECUTE_CODE",
|
||||
"nodeParams":[{"name":"value","type":"quote","quote":["b","output","value"],"quoteType":"output"}]}},
|
||||
{"id":"b","type":"code","properties":{"name":"B","action":"EXECUTE_CODE",
|
||||
"nodeParams":[{"name":"value","type":"quote","quote":["a","output","value"],"quoteType":"output"}]}}
|
||||
],
|
||||
"edges": []
|
||||
}
|
||||
""";
|
||||
|
||||
GlobalException error = assertThrows(GlobalException.class,
|
||||
() -> FlowModelBuilder.buildExecutableGraph(flow));
|
||||
|
||||
assertEquals("节点参数存在循环引用:A -> B -> A", error.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<TaskNodeExecuteContext> 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;
|
||||
}
|
||||
}
|
||||
|
||||
@ -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());
|
||||
|
||||
@ -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<String, FlowNodeWrapper> 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<FlowParamDef> startInputDefs,
|
||||
TaskNodeExecuteMessage rootMessage,
|
||||
Runnable onFinished,
|
||||
List<Integer> iterations) {
|
||||
executedNodeId = node.getNodeId();
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -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<String, FlowNodeWrapper> 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;
|
||||
}
|
||||
}
|
||||
@ -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;
|
||||
}
|
||||
}
|
||||
@ -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")));
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<FlowExecutionEvent> 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());
|
||||
}
|
||||
}
|
||||
@ -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<TaskNodeExecuteResult> 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();
|
||||
|
||||
@ -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> T mapper(Class<T> 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);
|
||||
}
|
||||
}
|
||||
@ -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());
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user