- 新增 SUB_FLOW 类型枚举值用于标识子流程节点 - 实现子流程异步媒体分析协调器管理机制 - 添加循环迭代路径追踪到执行上下文日志 - 增强流程控制服务支持嵌套子流程暂停恢复 - 实现子流程作用域内异步任务等待和取消功能 - 优化循环体多入口独立执行线程调度机制 - 添加节点参数循环引用检测验证机制 - 支持遗留循环参数值解析兼容性处理 - 配置虚拟机器人ID列表绕过在线检查机制
211 lines
9.3 KiB
Java
211 lines
9.3 KiB
Java
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;
|
||
import com.cmvr.inspection.mapper.InspectionRobotMapper;
|
||
import com.cmvr.quic.gateway.v1.ManagedDeviceStatus;
|
||
import com.cmvr.quic.gateway.v1.NodeEvent;
|
||
import com.cmvr.quic.gateway.v1.NodeEventType;
|
||
import com.cmvr.quic.gateway.v1.NodeSnapshot;
|
||
import lombok.RequiredArgsConstructor;
|
||
import org.apache.commons.lang3.StringUtils;
|
||
import org.slf4j.Logger;
|
||
import org.slf4j.LoggerFactory;
|
||
import org.springframework.beans.factory.annotation.Value;
|
||
import org.springframework.stereotype.Service;
|
||
import org.springframework.transaction.annotation.Transactional;
|
||
|
||
import java.util.ArrayList;
|
||
import java.util.Date;
|
||
import java.util.List;
|
||
import java.util.Objects;
|
||
|
||
@Service
|
||
@RequiredArgsConstructor
|
||
public class RobotQuicStateService {
|
||
|
||
private static final Logger log = LoggerFactory.getLogger(RobotQuicStateService.class);
|
||
private static final int AGV_DEVICE_KIND = 1;
|
||
|
||
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;
|
||
|
||
@Transactional(rollbackFor = Exception.class)
|
||
public void synchronize(NodeEvent event) {
|
||
NodeSnapshot node = event.getNode();
|
||
String robotId = StringUtils.trimToNull(node.getRobotId());
|
||
if (robotId == null) {
|
||
log.warn("忽略未携带机器人ID的QUIC事件,nodeId={},eventType={}",
|
||
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()) {
|
||
log.warn("QUIC机器人尚未在平台登记,robotId={},nodeId={}", robotId, node.getNodeId());
|
||
return;
|
||
}
|
||
if (matches.size() != 1) {
|
||
log.error("机器人ID存在{}条重复记录,拒绝同步,robotId={}", matches.size(), robotId);
|
||
return;
|
||
}
|
||
|
||
InspectionRobot robot = matches.get(0);
|
||
if ("1".equals(robot.getArchivedStatus())) {
|
||
log.warn("已归档机器人尝试连接,拒绝同步,robotId={}", robotId);
|
||
return;
|
||
}
|
||
|
||
boolean offline = event.getType() == NodeEventType.NODE_EVENT_TYPE_OFFLINE || !node.getOnline();
|
||
if (offline && StringUtils.isNotBlank(robot.getQuicSessionId())
|
||
&& !robot.getQuicSessionId().equals(node.getSessionId())) {
|
||
log.info("忽略旧QUIC会话的离线事件,robotId={},activeSession={},eventSession={}",
|
||
robotId, robot.getQuicSessionId(), node.getSessionId());
|
||
return;
|
||
}
|
||
|
||
Date eventTime = new Date(event.getEventTimeUnixMs() == 0
|
||
? System.currentTimeMillis() : event.getEventTimeUnixMs());
|
||
String oldHost = robot.getIpAddress();
|
||
Long oldPort = robot.getPort();
|
||
if (!offline) {
|
||
String host = resolveGrpcHost(node);
|
||
Long port = node.getGrpcEndpoint().getPort() == 0
|
||
? null : Long.valueOf(node.getGrpcEndpoint().getPort());
|
||
robot.setIpAddress(host);
|
||
robot.setPort(port);
|
||
}
|
||
|
||
robot.setConnectStatus(offline ? "0" : "1");
|
||
robot.setQuicNodeId(node.getNodeId());
|
||
robot.setQuicSessionId(node.getSessionId());
|
||
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");
|
||
}
|
||
robotMapper.updateQuicRuntimeState(robot);
|
||
|
||
if (offline) {
|
||
deviceMapper.markAllOffline(robotId, eventTime);
|
||
} else {
|
||
synchronizeDevices(robotId, node, eventTime);
|
||
}
|
||
|
||
if (offline || !Objects.equals(oldHost, robot.getIpAddress())
|
||
|| !Objects.equals(oldPort, robot.getPort())) {
|
||
grpcServiceManager.invalidateRobot(robotId);
|
||
}
|
||
if (event.getType() == NodeEventType.NODE_EVENT_TYPE_HEARTBEAT) {
|
||
log.debug("QUIC机器人心跳已同步,robotId={},nodeId={},grpc={}:{},devices={}",
|
||
robotId, node.getNodeId(), robot.getIpAddress(), robot.getPort(),
|
||
node.getDeviceManager().getDevicesCount());
|
||
} else {
|
||
log.info("QUIC机器人状态已同步,robotId={},nodeId={},online={},grpc={}:{},devices={}",
|
||
robotId, node.getNodeId(), !offline, robot.getIpAddress(), robot.getPort(),
|
||
node.getDeviceManager().getDevicesCount());
|
||
}
|
||
}
|
||
|
||
@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()) {
|
||
String deviceId = StringUtils.trimToNull(status.getDeviceId());
|
||
if (deviceId == null) {
|
||
continue;
|
||
}
|
||
reportedIds.add(deviceId);
|
||
List<InspectionRobotDevice> matches = deviceMapper.selectByIdentity(robotId, deviceId);
|
||
if (matches.size() > 1) {
|
||
log.error("机器人设备存在重复记录,跳过更新,robotId={},deviceId={},count={}",
|
||
robotId, deviceId, matches.size());
|
||
continue;
|
||
}
|
||
InspectionRobotDevice device = matches.isEmpty()
|
||
? createDevice(robotId, deviceId, eventTime) : matches.get(0);
|
||
device.setDeviceKind(status.getKind());
|
||
device.setTypeName(status.getTypeName());
|
||
device.setEnabled(status.getEnabled() ? "1" : "0");
|
||
device.setManagerState(status.getManagerState());
|
||
device.setHealthStatus(status.getHealth());
|
||
device.setOnlineStatus("1");
|
||
device.setHasError(status.getHasError() ? "1" : "0");
|
||
device.setErrorMessage(status.getErrorMessage());
|
||
device.setStatusUpdatedTime(status.getStatusUpdatedAtUnixMs() == 0
|
||
? null : new Date(status.getStatusUpdatedAtUnixMs()));
|
||
device.setLastSeenTime(eventTime);
|
||
device.setOfflineTime(null);
|
||
device.setUpdateBy("quic");
|
||
device.setUpdateTime(eventTime);
|
||
if (matches.isEmpty()) {
|
||
deviceMapper.insert(device);
|
||
} else {
|
||
deviceMapper.updateRuntimeState(device);
|
||
}
|
||
}
|
||
deviceMapper.markMissingOffline(robotId, reportedIds, eventTime);
|
||
}
|
||
|
||
private InspectionRobotDevice createDevice(String robotId, String deviceId, Date now) {
|
||
InspectionRobotDevice device = new InspectionRobotDevice();
|
||
device.setRobotId(robotId);
|
||
device.setDeviceId(deviceId);
|
||
device.setDeviceName(deviceId);
|
||
device.setCreateBy("quic");
|
||
device.setCreateTime(now);
|
||
return device;
|
||
}
|
||
|
||
private boolean hasEnabledAgv(NodeSnapshot node) {
|
||
return node.getDeviceManager().getDevicesList().stream()
|
||
.anyMatch(device -> device.getKind() == AGV_DEVICE_KIND && device.getEnabled());
|
||
}
|
||
|
||
private String resolveGrpcHost(NodeSnapshot node) {
|
||
String advertised = StringUtils.trimToNull(node.getGrpcEndpoint().getHost());
|
||
String observed = StringUtils.trimToNull(node.getObservedSourceIp());
|
||
if ("advertised".equalsIgnoreCase(grpcHostSource)) {
|
||
return advertised != null ? advertised : observed;
|
||
}
|
||
return observed != null ? observed : advertised;
|
||
}
|
||
}
|