"""设备状态检查服务""" import time from typing import List, Dict, Optional from sqlalchemy.orm import Session from app.services.ssh_service import SSHService from app.models.device import OLTDevice, ONUDevice, DeviceStatusHistory, DuplicateMac, NewDevice from datetime import datetime import re # SSH 连接池缓存,TTL 5 分钟 _conn_pool: Dict[int, tuple[SSHService, float]] = {} _POOL_TTL = 300 def _get_cached_ssh(olt_ip: str, olt_user: str, olt_pass: str, olt_id: int) -> SSHService: """获取缓存的 SSH 连接,过期自动重连""" entry = _conn_pool.get(olt_id) if entry: ssh, ts = entry if time.time() - ts < _POOL_TTL: return ssh try: ssh.close() except Exception: pass ssh = SSHService(olt_ip, olt_user, olt_pass) ssh.connect() _conn_pool[olt_id] = (ssh, time.time()) return ssh def parse_distance(distance_str: Optional[str]) -> Optional[int]: """将距离字符串转为整数,如 '<1000' -> 1000, '1234' -> 1234""" if not distance_str: return None m = re.search(r'\d+', distance_str) return int(m.group()) if m else None class CheckService: """设备状态检查服务""" def __init__(self, db: Session): self.db = db def update_status_only(self, olt_id: int) -> Dict: """扫描 OLT,仅更新已有设备的在线状态和距离(按全局 MAC 匹配)。 不修改 olt_id/端口等字段,不标记 unknown,不入库新设备。 """ olt = self.db.query(OLTDevice).filter(OLTDevice.id == olt_id).first() if not olt: raise Exception(f"OLT 设备不存在: {olt_id}") ssh = _get_cached_ssh(olt.ip_address, olt.username, olt.password, olt_id) try: output = ssh.execute_command(olt.slot_command) onu_info_dict, _ = ssh.parse_onu_info(output) checked_at = datetime.utcnow() online_count = 0 offline_count = 0 # 全局 MAC → (id, model) 索引 existing = { row.mac_address.lower(): (row.id, row.model) for row in self.db.query(ONUDevice.mac_address, ONUDevice.id, ONUDevice.model).all() } for mac, onu_info in onu_info_dict.items(): entry = existing.get(mac) if entry is None: continue # 不在库中,跳过 onu_id, onu_model = entry status = onu_info.status distance_m = parse_distance(onu_info.distance_str) if status == 'online': online_count += 1 else: offline_count += 1 # 有值时同步端口和 OLT 归属(不用 None 覆盖现有值) update_fields = {"olt_id": olt_id} if onu_info.slot_number is not None: update_fields["slot_number"] = onu_info.slot_number if onu_info.port_number is not None: update_fields["port_number"] = onu_info.port_number if onu_info.port_id: update_fields["port_id"] = onu_info.port_id if onu_info.model and not onu_model: update_fields["model"] = onu_info.model self.db.query(ONUDevice).filter(ONUDevice.id == onu_id).update(update_fields) self.db.add(DeviceStatusHistory( onu_device_id=onu_id, status=status, distance_m=distance_m, checked_at=checked_at, response_data=None, )) self.db.commit() return { "online": online_count, "offline": offline_count, } except Exception: # 连接异常时清除缓存,下次自动重连 _conn_pool.pop(olt_id, None) raise def check_single_device(self, device_id: int) -> Dict: """通过 SSH 单独查询一台 ONU 设备的当前状态和距离。 命令格式: display onu slot {slot} | include {mac} """ device = self.db.query(ONUDevice).filter(ONUDevice.id == device_id).first() if not device: raise Exception("设备不存在") if not device.olt_id: raise Exception("该设备未关联 OLT,无法查询") olt = self.db.query(OLTDevice).filter(OLTDevice.id == device.olt_id).first() if not olt: raise Exception("关联的 OLT 不存在") mac = device.mac_address.lower() cmd = f"{olt.slot_command} | include {mac}" ssh = SSHService(olt.ip_address, olt.username, olt.password) try: ssh.connect() output = ssh.execute_command(cmd) # 跳过命令回显行(含 MAC 但无端口标识),只接受有 port_id 的行 info = None for line in output.splitlines(): if mac in line.lower(): parsed = ssh._parse_device_line(line, device.slot_number) if parsed and parsed.port_id: info = parsed break if info is None: status = 'offline' distance_m = None else: status = info.status distance_m = parse_distance(info.distance_str) if info.model and not device.model: device.model = info.model # 同步端口信息到数据库 device.slot_number = info.slot_number device.port_number = info.port_number device.port_id = info.port_id self.db.add(DeviceStatusHistory( onu_device_id=device.id, status=status, distance_m=distance_m, checked_at=datetime.utcnow(), response_data=None, )) self.db.commit() olt_location = olt.location if olt else None return { "status": status, "distance_m": distance_m, "model": device.model, "port_id": device.port_id, "slot_number": device.slot_number, "port_number": device.port_number, "olt_location": olt_location, } finally: ssh.close() def scan_and_discover(self, olt_id: int) -> Dict: """扫描 OLT,更新已有设备状态,并将新发现的在线设备入库。 不标记 unknown,不修改已有设备的 olt_id/端口以外的字段。 """ olt = self.db.query(OLTDevice).filter(OLTDevice.id == olt_id).first() if not olt: raise Exception(f"OLT 设备不存在: {olt_id}") ssh = _get_cached_ssh(olt.ip_address, olt.username, olt.password, olt_id) try: output = ssh.execute_command(olt.slot_command) onu_info_dict, duplicate_dict = ssh.parse_onu_info(output) checked_at = datetime.utcnow() online_count = 0 offline_count = 0 new_count = 0 existing = { row.mac_address.lower(): (row.id, row.model) for row in self.db.query(ONUDevice.mac_address, ONUDevice.id, ONUDevice.model).all() } for mac, onu_info in onu_info_dict.items(): entry = existing.get(mac) if entry is None: # 新设备:只入库在线的 if onu_info.status != 'online': continue onu = ONUDevice( mac_address=mac, olt_id=olt_id, slot_number=onu_info.slot_number, port_number=onu_info.port_number, port_id=onu_info.port_id, distance_m=parse_distance(onu_info.distance_str), loid=onu_info.loid, model=onu_info.model, ) self.db.add(onu) self.db.flush() self.db.add(NewDevice(onu_device_id=onu.id, olt_id=olt_id)) new_count += 1 onu_id = onu.id else: onu_id, onu_model = entry # 同步端口和 OLT 归属(扫描结果以当前 OLT 为准) self.db.query(ONUDevice).filter(ONUDevice.id == onu_id).update({ "olt_id": olt_id, "slot_number": onu_info.slot_number, "port_number": onu_info.port_number, "port_id": onu_info.port_id, }) if onu_info.model and not onu_model: self.db.query(ONUDevice).filter(ONUDevice.id == onu_id).update({"model": onu_info.model}) status = onu_info.status distance_m = parse_distance(onu_info.distance_str) if status == 'online': online_count += 1 else: offline_count += 1 self.db.add(DeviceStatusHistory( onu_device_id=onu_id, status=status, distance_m=distance_m, checked_at=checked_at, response_data=None, )) self._save_duplicate_macs(olt_id, duplicate_dict) self.db.commit() return { "online": online_count, "offline": offline_count, "new_discovered": new_count, } except Exception: _conn_pool.pop(olt_id, None) raise async def scan_olt(self, olt_id: int) -> Dict: """仅扫描 OLT,返回发现的设备列表(不写入数据库)""" olt = self.db.query(OLTDevice).filter(OLTDevice.id == olt_id).first() if not olt: raise Exception(f"OLT 设备不存在: {olt_id}") ssh = _get_cached_ssh(olt.ip_address, olt.username, olt.password, olt_id) try: output = ssh.execute_command(olt.slot_command) onu_info_dict, duplicate_dict = ssh.parse_onu_info(output) # 对比全局 MAC existing_macs = { onu.mac_address.lower() for onu in self.db.query(ONUDevice.mac_address).all() } devices = [] for mac, info in onu_info_dict.items(): devices.append({ "mac_address": mac, "status": info.status, "distance_m": info.distance_str, "slot_number": info.slot_number, "port_number": info.port_number, "port_id": info.port_id, "loid": info.loid, "model": info.model, "is_new": mac not in existing_macs, }) duplicates = [] for mac, records in duplicate_dict.items(): duplicates.append({ "mac_address": mac, "ports": [{"port_id": r.port_id, "status": r.status} for r in records], }) return { "olt_id": olt_id, "olt_ip": olt.ip_address, "total": len(devices), "new": sum(1 for d in devices if d["is_new"]), "devices": devices, "duplicates": duplicates, } except Exception: _conn_pool.pop(olt_id, None) raise def _save_duplicate_macs(self, olt_id: int, duplicate_dict: dict): for mac, records in duplicate_dict.items(): ports = [{"port_id": r.port_id, "status": r.status} for r in records] existing = self.db.query(DuplicateMac).filter( DuplicateMac.olt_id == olt_id, DuplicateMac.mac_address == mac ).first() if existing: existing.ports = ports existing.last_seen_at = datetime.utcnow() else: self.db.add(DuplicateMac(olt_id=olt_id, mac_address=mac, ports=ports))