package com.cmvr.resource.service; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.cmvr.common.utils.DateUtils; import com.cmvr.edge.client.manage.GrpcServiceManager; 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 com.cmvr.resource.config.ResourceRobotStateProperties; import com.cmvr.resource.domain.ResourceDevice; import com.cmvr.resource.domain.ResourceRobot; import com.cmvr.resource.domain.ResourceRobotTerminal; import com.cmvr.resource.domain.ResourceTerminal; import com.cmvr.resource.mapper.ResourceDeviceMapper; import com.cmvr.resource.mapper.ResourceRobotMapper; import com.cmvr.resource.mapper.ResourceRobotTerminalMapper; import com.cmvr.resource.mapper.ResourceTerminalMapper; 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; import java.util.UUID; @Service @RequiredArgsConstructor public class ResourceQuicStateSynchronizer { private static final Logger log = LoggerFactory.getLogger(ResourceQuicStateSynchronizer.class); private static final int AGV_DEVICE_KIND = 1; private final ResourceRobotMapper robotMapper; private final ResourceDeviceMapper deviceMapper; private final ResourceTerminalMapper terminalMapper; private final ResourceRobotTerminalMapper robotTerminalMapper; private final ResourceRobotStateProperties stateProperties; private final GrpcServiceManager grpcServiceManager; @Value("${cmvr.quic.grpc-host-source:observed-source}") private String grpcHostSource; @Transactional(rollbackFor = Exception.class) public boolean synchronize(NodeEvent event) { NodeSnapshot node = event.getNode(); String robotId = StringUtils.trimToNull(node.getRobotId()); if (robotId == null || stateProperties.isVirtualRobot(robotId)) { return false; } ResourceRobot robot = robotMapper.selectByRobotIdentity(robotId); if (robot == null || "1".equals(robot.getArchivedStatus())) { log.warn("QUIC robot is not registered in resource center, robotId={}, nodeId={}", robotId, node.getNodeId()); return false; } boolean offline = event.getType() == NodeEventType.NODE_EVENT_TYPE_OFFLINE || !node.getOnline(); if (offline && StringUtils.isNotBlank(robot.getQuicSessionId()) && !robot.getQuicSessionId().equals(node.getSessionId())) { return true; } String oldHost = robot.getIpAddress(); Long oldPort = robot.getPort(); Date eventTime = new Date(event.getEventTimeUnixMs() == 0 ? System.currentTimeMillis() : event.getEventTimeUnixMs()); if (!offline) { robot.setIpAddress(resolveGrpcHost(node)); robot.setPort(node.getGrpcEndpoint().getPort() == 0 ? null : (long) node.getGrpcEndpoint().getPort()); } 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); } if (offline || !hasEnabledAgv(node)) { robot.setStatus("1"); } robot.setUpdateTime(DateUtils.getNowDate()); robotMapper.updateQuicRuntimeState(robot); synchronizeTerminal(robotId, node, eventTime, offline); 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); } return true; } @Transactional(rollbackFor = Exception.class) public int expireStaleRobots() { Date now = DateUtils.getNowDate(); Date cutoff = stateProperties.heartbeatCutoff(now); int count = 0; for (ResourceRobot robot : robotMapper.selectStaleConnected( cutoff, stateProperties.normalizedVirtualRobotIds())) { if (robotMapper.markHeartbeatExpiredOffline(robot.getId(), cutoff, now) == 0) { continue; } deviceMapper.markAllOffline(robot.getRobotId(), now); grpcServiceManager.invalidateRobot(robot.getRobotId()); count++; } return count; } private void synchronizeDevices(String robotId, NodeSnapshot node, Date eventTime) { List reportedIds = new ArrayList<>(); for (ManagedDeviceStatus status : node.getDeviceManager().getDevicesList()) { String deviceId = StringUtils.trimToNull(status.getDeviceId()); if (deviceId == null) { continue; } reportedIds.add(deviceId); List matches = deviceMapper.selectByIdentity(robotId, deviceId); if (matches.size() > 1) { log.error("Duplicate resource device identity, robotId={}, deviceId={}", robotId, deviceId); continue; } ResourceDevice 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 ResourceDevice createDevice(String robotId, String deviceId, Date now) { ResourceDevice device = new ResourceDevice(); device.setId(UUID.randomUUID().toString().replace("-", "")); device.setRobotId(robotId); device.setDeviceId(deviceId); device.setDeviceName(deviceId); device.setCreateBy("quic"); device.setCreateTime(now); return device; } private void synchronizeTerminal(String robotId, NodeSnapshot node, Date eventTime, boolean offline) { String terminalCode = StringUtils.trimToNull(node.getNodeId()); if (terminalCode == null) { return; } ResourceTerminal terminal = terminalMapper.selectByCode(terminalCode); boolean newTerminal = terminal == null; if (newTerminal) { terminal = new ResourceTerminal(); terminal.setId(UUID.randomUUID().toString().replace("-", "")); terminal.setTerminalCode(terminalCode); terminal.setTerminalName(terminalCode); terminal.setCreateBy("quic"); terminal.setCreateTime(eventTime); } terminal.setIpAddress(StringUtils.trimToNull(node.getObservedSourceIp())); terminal.setPort(node.getGrpcEndpoint().getPort() == 0 ? null : (long) node.getGrpcEndpoint().getPort()); terminal.setSoftwareVersion(node.getSoftwareVersion()); terminal.setConnectStatus(offline ? "0" : "1"); terminal.setLastHeartbeatTime(eventTime); terminal.setStatus(offline ? "OFFLINE" : "ONLINE"); terminal.setUpdateBy("quic"); terminal.setUpdateTime(eventTime); if (newTerminal) { terminalMapper.insert(terminal); } else { terminalMapper.updateById(terminal); } ResourceRobotTerminal binding = robotTerminalMapper.selectOne( new LambdaQueryWrapper() .eq(ResourceRobotTerminal::getRobotId, robotId) .eq(ResourceRobotTerminal::getTerminalId, terminal.getId())); if (binding == null) { binding = new ResourceRobotTerminal(); binding.setId(UUID.randomUUID().toString().replace("-", "")); binding.setRobotId(robotId); binding.setTerminalId(terminal.getId()); binding.setBindTime(eventTime); binding.setCreateTime(eventTime); binding.setActive(offline ? "0" : "1"); robotTerminalMapper.insert(binding); } else { binding.setActive(offline ? "0" : "1"); binding.setUnbindTime(offline ? eventTime : null); binding.setUpdateTime(eventTime); robotTerminalMapper.updateById(binding); } if (!offline) { robotTerminalMapper.deactivateOtherBindings(robotId, terminal.getId()); } } 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; } }