From 0d85b87be79723d6688bbd3c1167faeab2fb4a47 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 31 Jul 2026 21:30:46 +0800 Subject: [PATCH] Add Huawei/Cisco port traffic monitoring and keep Cisco CLI on send_command. Extend port discovery/sampling beyond ZTE, harden IOSv session prep, and avoid send_command_timing so Building configuration does not truncate config sync. Co-authored-by: Cursor --- netx_api/config_sync_runner.py | 3 +- netx_api/ne_collect_runner.py | 3 +- netx_api/ne_connect.py | 3 +- netx_api/ne_netmiko.py | 19 ++ netx_api/ne_session_factory.py | 39 +++- netx_api/port_traffic_commands.py | 14 +- netx_api/port_traffic_parsers.py | 298 ++++++++++++++++++++++++++++- netx_api/port_traffic_runner.py | 18 +- netx_api/port_traffic_service.py | 16 +- tests/test_port_traffic_parsers.py | 156 ++++++++++++++- web/WEB.md | 2 +- 11 files changed, 544 insertions(+), 27 deletions(-) diff --git a/netx_api/config_sync_runner.py b/netx_api/config_sync_runner.py index ffdb969..503968a 100644 --- a/netx_api/config_sync_runner.py +++ b/netx_api/config_sync_runner.py @@ -21,6 +21,7 @@ from .config_sync_service import finalize_cycle, sync_cycle_progress from .db import SessionLocal from .models import ConfigSyncCycle, ConfigSyncPolicy, ConfigSyncTask, NeConfigHistory, NeConfigSnapshot from .ne_cli_errors import format_cli_failure, session_log_text +from .ne_netmiko import send_show_command from .ne_session_factory import close_netmiko_connection, open_netmiko_connection _log = logging.getLogger("netx.config_sync.runner") @@ -113,7 +114,7 @@ def _collect_commands(creds: dict[str, Any], commands: list[str]) -> list[str]: outputs: list[str] = [] for command in commands: try: - out = conn.send_command(command_string=command, read_timeout=per_cmd) + out = send_show_command(conn, command, read_timeout=per_cmd) except Exception as exc: raise RuntimeError(format_cli_failure(exc, session_log_text(log_buf))) from exc outputs.append(str(out or "")) diff --git a/netx_api/ne_collect_runner.py b/netx_api/ne_collect_runner.py index f90cc72..11f3cbf 100644 --- a/netx_api/ne_collect_runner.py +++ b/netx_api/ne_collect_runner.py @@ -15,6 +15,7 @@ from .db import SessionLocal from .models import ManagedNE, NeCollectionJob, NeCollectionRun from .ne_collection_paths import clear_run_output_files, run_output_dir from .ne_crypto import CredentialCryptoError +from .ne_netmiko import send_show_command from .ne_service import get_device_credentials from .ne_session_factory import close_netmiko_connection, open_netmiko_connection @@ -52,7 +53,7 @@ def _collect_on_device(creds: dict[str, Any], commands: list[str]) -> str: for command in commands: ts = datetime.now().isoformat(timespec="seconds") chunks.append(f'>>> [{ts}] {{"String":"{command}", "Match":"{prompt}", "Timeout":0}}\n') - out = conn.send_command(command_string=command, read_timeout=per_cmd) + out = send_show_command(conn, command, read_timeout=per_cmd) chunks.append(str(out or "")) chunks.append("\n") return "".join(chunks) diff --git a/netx_api/ne_connect.py b/netx_api/ne_connect.py index 309eed8..ff9a107 100644 --- a/netx_api/ne_connect.py +++ b/netx_api/ne_connect.py @@ -13,6 +13,7 @@ from .models import ManagedNE, UmeCliOverride, UmeInventoryNE from .ne_crypto import CredentialCryptoError from .cli_resolve import resolve_cli_target from .ne_service import get_device_credentials +from .ne_netmiko import send_show_command from .ne_session_factory import ( bastion_ssh_cli, close_netmiko_connection, @@ -238,7 +239,7 @@ def _probe_device(creds: dict[str, Any]) -> tuple[str, str, str | None, str]: command = hostname_probe_command(creds["device_type"], vendor) output = "" if command: - output = conn.send_command(command_string=command, read_timeout=_PROBE_READ_TIMEOUT) + output = send_show_command(conn, command, read_timeout=_PROBE_READ_TIMEOUT) hostname = parse_hostname_from_output(creds["device_type"], vendor, output, prompt) if hostname: msg = f"connected: {hostname}" diff --git a/netx_api/ne_netmiko.py b/netx_api/ne_netmiko.py index ec6b97b..e60620a 100644 --- a/netx_api/ne_netmiko.py +++ b/netx_api/ne_netmiko.py @@ -2,6 +2,8 @@ from __future__ import annotations +from typing import Any + def normalize_netmiko_device_type(device_type: str, protocol: str) -> str: dt = str(device_type or "").strip() @@ -15,3 +17,20 @@ def normalize_netmiko_device_type(device_type: str, protocol: str) -> str: if "telnet" not in dt and "ssh" not in dt: return f"{dt}_{proto}" return dt + + +def is_cisco_ios_device_type(device_type: str) -> bool: + dt = str(device_type or "").strip().lower() + return "cisco_ios" in dt or dt in {"cisco_ios", "cisco_ios_ssh", "cisco_ios_telnet"} + + +def send_show_command(conn: Any, command: str, *, read_timeout: int = 120) -> str: + """Send a show/display command via ``send_command`` (wait for device prompt). + + Do not use ``send_command_timing`` for Cisco config collection: long idle during + ``Building configuration...`` is treated as end-of-output and truncates the config. + """ + cmd = str(command or "").strip() + if not cmd: + return "" + return str(conn.send_command(command_string=cmd, read_timeout=read_timeout) or "") diff --git a/netx_api/ne_session_factory.py b/netx_api/ne_session_factory.py index 392d3bb..ee4d1b5 100644 --- a/netx_api/ne_session_factory.py +++ b/netx_api/ne_session_factory.py @@ -259,13 +259,44 @@ def _interactive_driver_class(base_cls: type) -> type: return _InteractiveSession +def _cisco_ios_collection_driver_class(base_cls: type) -> type: + """Cisco IOSv-friendly session prep: avoid cmd_verify on terminal width/length.""" + + class _CiscoIosCollectionSession(base_cls): # type: ignore[misc,valid-type] + def session_preparation(self) -> None: + # Default Netmiko waits for exact echo of "terminal width 511" (ReadTimeout on IOSv). + self._test_channel_read(pattern=r"[>#]") + try: + self.set_terminal_width(command="terminal width 511", pattern=r"[>#]") + except Exception: + pass + try: + self.disable_paging(command="terminal length 0", cmd_verify=False, pattern=r"[>#]") + except Exception: + try: + self.send_command_timing("terminal length 0", read_timeout=15) + except Exception: + pass + self.set_base_prompt() + + _CiscoIosCollectionSession.__name__ = ( + f"CiscoIosCollection{getattr(base_cls, '__name__', 'Netmiko')}" + ) + return _CiscoIosCollectionSession + + def _build_netmiko_connection(dev: dict[str, Any], *, interactive: bool = False) -> ConnectHandler: """Instantiate Netmiko from connect kwargs; optional interactive skips paging cmds.""" - if not interactive: - return ConnectHandler(**dev) + from .ne_netmiko import is_cisco_ios_device_type + device_type = str(dev.get("device_type") or "").strip() - base_cls = _netmiko_driver_class(device_type) - return _interactive_driver_class(base_cls)(**dev) + if interactive: + base_cls = _netmiko_driver_class(device_type) + return _interactive_driver_class(base_cls)(**dev) + if is_cisco_ios_device_type(device_type): + base_cls = _netmiko_driver_class(device_type) + return _cisco_ios_collection_driver_class(base_cls)(**dev) + return ConnectHandler(**dev) def _netmiko_over_ssh_client( diff --git a/netx_api/port_traffic_commands.py b/netx_api/port_traffic_commands.py index 1f6b681..b747367 100644 --- a/netx_api/port_traffic_commands.py +++ b/netx_api/port_traffic_commands.py @@ -1,4 +1,4 @@ -"""Vendor → port traffic CLI command matrix (ZTE first).""" +"""Vendor → port traffic CLI command matrix (ZTE / Huawei / Cisco).""" from __future__ import annotations @@ -22,6 +22,18 @@ def commands_for_vendor(vendor: str, device_type: str = "") -> PortTrafficComman detail_template="show interface {ifname}", vendor_key=key, ) + if key == "huawei": + return PortTrafficCommands( + brief="display interface brief", + detail_template="display interface {ifname}", + vendor_key=key, + ) + if key == "cisco": + return PortTrafficCommands( + brief="show ip interface brief", + detail_template="show interfaces {ifname}", + vendor_key=key, + ) return None diff --git a/netx_api/port_traffic_parsers.py b/netx_api/port_traffic_parsers.py index 02d8603..fec42d9 100644 --- a/netx_api/port_traffic_parsers.py +++ b/netx_api/port_traffic_parsers.py @@ -1,4 +1,4 @@ -"""Parsers for ZTE show interface brief / detail (rate bit/s).""" +"""Parsers for ZTE / Huawei / Cisco interface brief & detail (rate bit/s).""" from __future__ import annotations @@ -15,8 +15,9 @@ _BW_UNIT = { } _RE_BW_COMPACT = re.compile(r"^(\d+(?:\.\d+)?)\s*([kKmMgGtT])(?:bit)?s?$", re.I) +# ZTE "BW 1 Gbit/s" and Cisco "BW 1000000 Kbit/sec" _RE_BW_DETAIL = re.compile( - r"\bBW\s+(\d+(?:\.\d+)?)\s*([kKmMgGtT])?\s*(?:G?bit|bit)/s\b", + r"\bBW\s+(\d+(?:\.\d+)?)\s*([kKmMgGtT])?\s*(?:G?bit|bit)/s(?:ec)?\b", re.I, ) _RE_RATE_PERIOD = re.compile(r"Rate\s+period\s*:\s*(\d+)\s*s", re.I) @@ -30,14 +31,53 @@ _RE_IF_UP = re.compile(r"^(\S+)\s+is\s+(up|down)\b", re.I | re.M) _RE_DESC = re.compile(r"^\s*Description:\s*(.+?)\s*$", re.I | re.M) _RE_PROMPT_LINE = re.compile(r"[#>]\s*$") _RE_UPDOWN = re.compile(r"^(up|down)$", re.I) -# ZTE / common logical iface names seen on ZXROS (plus numbered variants). +_RE_UPDOWN_TOKEN = re.compile(r"^(up|down)\b", re.I) +# ZTE / Huawei / Cisco / common logical iface names. _RE_IFNAME = re.compile( r"^(?:" r"xxvgei|xgei|cgei|gei|fei|qli|smartgroup|bvi|vlan|loopback|mgmt|" - r"null|pos|atm|tunnel|irb|pw|eth|ethernet|port-channel|bundle" + r"null|pos|atm|tunnel|irb|pw|eth|ethernet|port-channel|bundle|" + r"gigabitethernet|xgigabitethernet|fastethernet|tengigabitethernet|" + r"hundredgige|fivegige|fortygige|serial|dialer|cellular|multilink|" + r"10ge|25ge|40ge|100ge|eth-trunk|vlanif|meth|loopback" r")[\w./:-]*$", re.I, ) +_RE_CISCO_BRIEF_ROW = re.compile( + r"^(\S+)\s+(\S+)\s+(YES|NO)\s+(\S+)\s+(.+?)\s+(up|down)\s*$", + re.I, +) +_RE_CISCO_IF_STATE = re.compile( + r"^(\S+)\s+is\s+(administratively\s+)?(up|down),\s*line\s+protocol\s+is\s+(up|down)\b", + re.I | re.M, +) +_RE_CISCO_RATE = re.compile( + r"(\d+)\s+(second|minute|hour)s?\s+(input|output)\s+rate\s+([\d.]+)\s*bits/sec", + re.I, +) +_RE_HW_BRIEF_ROW = re.compile( + r"^(\S+)\s+(\S+)\s+(\S+)\s+(\S+)\s+(\S+)\s+(\d+)\s+(\d+)\s*$" +) +_RE_HW_STATE = re.compile( + r"^(\S+)\s+current\s+state\s*:\s*(UP|DOWN|Administratively\s+DOWN)\b", + re.I | re.M, +) +_RE_HW_IN_RATE = re.compile( + r"Last\s+(\d+)\s+seconds\s+input\s+rate\s*:\s*([\d.]+)\s*bits/sec", + re.I, +) +_RE_HW_OUT_RATE = re.compile( + r"Last\s+(\d+)\s+seconds\s+output\s+rate\s*:\s*([\d.]+)\s*bits/sec", + re.I, +) +_RE_HW_IN_UTIL = re.compile( + r"Last\s+(\d+)\s+seconds\s+input\s+utility\s+rate\s*:\s*([\d.]+)\s*%", + re.I, +) +_RE_HW_OUT_UTIL = re.compile( + r"Last\s+(\d+)\s+seconds\s+output\s+utility\s+rate\s*:\s*([\d.]+)\s*%", + re.I, +) # Fixed-width columns from ZTE `show interface brief` header. _COL_IF = (0, 24) @@ -277,3 +317,253 @@ def detail_to_dict(row: DetailRates) -> dict[str, Any]: "in_util_pct": row.in_util_pct, "out_util_pct": row.out_util_pct, } + + +def _norm_updown(token: str) -> str: + """Normalize Huawei tokens like up(s) / *down to up|down|''.""" + text = (token or "").strip().lower() + if text.startswith("*"): + text = text[1:] + m = _RE_UPDOWN_TOKEN.match(text) + return m.group(1).lower() if m else "" + + +def parse_huawei_interface_brief(text: str) -> list[BriefPort]: + """Parse Huawei `display interface brief` into port rows (BW ignored / 0).""" + lines = (text or "").replace("\r\n", "\n").replace("\r", "\n").split("\n") + out: list[BriefPort] = [] + started = False + for raw in lines: + line = raw.rstrip() + if not line.strip(): + continue + low = line.lstrip().lower() + if low.startswith("interface") and "phy" in low and ("inuti" in low or "protocol" in low): + started = True + continue + if not started: + continue + stripped = line.strip() + if stripped.startswith("<") and stripped.endswith(">"): + break + if "#" in stripped and " " not in stripped.split("#", 1)[0]: + break + m = _RE_HW_BRIEF_ROW.match(stripped) + if not m: + continue + ifname = m.group(1) + phy = _norm_updown(m.group(2)) + prot = _norm_updown(m.group(3)) + if not ifname or ifname.lower() == "interface": + continue + if "#" in ifname or ">" in ifname: + continue + if not _RE_IFNAME.match(ifname): + continue + if not phy or not prot: + continue + out.append( + BriefPort( + ifname=ifname, + attribute="", + mode="", + bw_raw="", + bw_bps=0, + admin=phy, # Huawei brief has no separate Admin; PHY is closest. + phy=phy, + prot=prot, + description="", + ) + ) + return out + + +def parse_huawei_interface_detail(text: str) -> DetailRates: + """Parse Huawei `display interface {if}` Last N seconds rate / utility.""" + blob = text or "" + ifname = "" + admin_oper = "" + m_state = _RE_HW_STATE.search(blob) + if m_state: + ifname = m_state.group(1) + st = m_state.group(2).lower() + admin_oper = "down" if "down" in st else "up" + desc = "" + m_desc = _RE_DESC.search(blob) + if m_desc: + desc = m_desc.group(1).strip() + + rate_period = 0 + in_bps = 0.0 + out_bps = 0.0 + m_in = _RE_HW_IN_RATE.search(blob) + if m_in: + rate_period = int(m_in.group(1)) + in_bps = float(m_in.group(2)) + m_out = _RE_HW_OUT_RATE.search(blob) + if m_out: + if not rate_period: + rate_period = int(m_out.group(1)) + out_bps = float(m_out.group(2)) + + in_util = 0.0 + out_util = 0.0 + m_iu = _RE_HW_IN_UTIL.search(blob) + if m_iu: + if not rate_period: + rate_period = int(m_iu.group(1)) + in_util = float(m_iu.group(2)) + m_ou = _RE_HW_OUT_UTIL.search(blob) + if m_ou: + if not rate_period: + rate_period = int(m_ou.group(1)) + out_util = float(m_ou.group(2)) + + return DetailRates( + ifname=ifname, + admin_oper=admin_oper, + description=desc, + bw_bps=0, # sample has no BW; leave 0 for now + rate_period_sec=rate_period, + in_bps=in_bps, + out_bps=out_bps, + in_util_pct=in_util, + out_util_pct=out_util, + ) + + +def parse_interface_brief(text: str, vendor_key: str = "zte") -> list[BriefPort]: + key = str(vendor_key or "zte").strip().lower() + if key == "huawei": + return parse_huawei_interface_brief(text) + if key == "cisco": + return parse_cisco_interface_brief(text) + return parse_zte_interface_brief(text) + + +def parse_interface_detail(text: str, vendor_key: str = "zte") -> DetailRates: + key = str(vendor_key or "zte").strip().lower() + if key == "huawei": + return parse_huawei_interface_detail(text) + if key == "cisco": + return parse_cisco_interface_detail(text) + return parse_zte_interface_detail(text) + + +def _cisco_period_to_sec(n: int, unit: str) -> int: + u = (unit or "").lower() + if u.startswith("second"): + return int(n) + if u.startswith("minute"): + return int(n) * 60 + if u.startswith("hour"): + return int(n) * 3600 + return int(n) + + +def parse_cisco_interface_brief(text: str) -> list[BriefPort]: + """Parse Cisco `show ip interface brief` into port rows.""" + lines = (text or "").replace("\r\n", "\n").replace("\r", "\n").split("\n") + out: list[BriefPort] = [] + started = False + for raw in lines: + line = raw.rstrip() + if not line.strip(): + continue + low = line.lstrip().lower() + if low.startswith("interface") and "status" in low and "protocol" in low: + started = True + continue + if not started: + continue + stripped = line.strip() + if stripped.endswith("#") or (stripped.endswith(">") and stripped.startswith("<")): + break + if "#" in stripped and " " not in stripped.split("#", 1)[0]: + break + m = _RE_CISCO_BRIEF_ROW.match(stripped) + if not m: + continue + ifname = m.group(1) + status = (m.group(5) or "").strip().lower() + prot = (m.group(6) or "").strip().lower() + if not _RE_IFNAME.match(ifname): + continue + if prot not in ("up", "down"): + continue + admin_down = "administratively" in status + if "down" in status: + phy = "down" + elif "up" in status: + phy = "up" + else: + continue + admin = "down" if admin_down else phy + out.append( + BriefPort( + ifname=ifname, + attribute="", + mode="", + bw_raw="", + bw_bps=0, + admin=admin, + phy=phy, + prot=prot, + description="", + ) + ) + return out + + +def parse_cisco_interface_detail(text: str) -> DetailRates: + """Parse Cisco `show interfaces {if}` BW + N minute/second rate (util ignored).""" + blob = text or "" + ifname = "" + admin_oper = "" + m_state = _RE_CISCO_IF_STATE.search(blob) + if m_state: + ifname = m_state.group(1) + admin_oper = "down" if m_state.group(2) or m_state.group(3).lower() == "down" else "up" + if not ifname: + m_up = _RE_IF_UP.search(blob) + if m_up: + ifname = m_up.group(1) + admin_oper = m_up.group(2).lower() + + desc = "" + m_desc = _RE_DESC.search(blob) + if m_desc: + desc = m_desc.group(1).strip() + + bw_bps = 0 + for line in blob.splitlines(): + if "BW" in line.upper() and ("bit" in line.lower()): + bw_bps = parse_bw_to_bps(line) + if bw_bps: + break + + rate_period = 0 + in_bps = 0.0 + out_bps = 0.0 + for m in _RE_CISCO_RATE.finditer(blob): + period = _cisco_period_to_sec(int(m.group(1)), m.group(2)) + direction = m.group(3).lower() + rate = float(m.group(4)) + if not rate_period: + rate_period = period + if direction == "input": + in_bps = rate + else: + out_bps = rate + + return DetailRates( + ifname=ifname, + admin_oper=admin_oper, + description=desc, + bw_bps=bw_bps, + rate_period_sec=rate_period, + in_bps=in_bps, + out_bps=out_bps, + in_util_pct=0.0, # sample has no util %; ignore for now + out_util_pct=0.0, + ) diff --git a/netx_api/port_traffic_runner.py b/netx_api/port_traffic_runner.py index eb06d26..dc78b38 100644 --- a/netx_api/port_traffic_runner.py +++ b/netx_api/port_traffic_runner.py @@ -17,8 +17,9 @@ from .config import settings from .db import SessionLocal from .models import PortTrafficSample, PortTrafficTarget, PortTrafficTask from .ne_session_factory import close_netmiko_connection, open_netmiko_connection +from .ne_netmiko import send_show_command from .port_traffic_commands import commands_for_vendor, detail_command -from .port_traffic_parsers import parse_zte_interface_detail +from .port_traffic_parsers import parse_interface_detail _log = logging.getLogger("netx.port_traffic.runner") _pools: dict[str, ThreadPoolExecutor] = {} @@ -129,7 +130,7 @@ def _finish_collect_round(task_id: str, *, error: str = "") -> None: def _run_show(creds: dict[str, Any], command: str, read_timeout: int) -> str: conn = open_netmiko_connection(creds, session_timeout=read_timeout + 60) try: - return str(conn.send_command(command_string=command, read_timeout=read_timeout) or "") + return send_show_command(conn, command, read_timeout=read_timeout) finally: close_netmiko_connection(conn) @@ -137,6 +138,7 @@ def _run_show(creds: dict[str, Any], command: str, read_timeout: int) -> str: def _sample_one_target(target_row_id: str) -> None: creds: dict[str, Any] | None = None cmd = "" + vendor_key = "zte" per_cmd = int(settings.ne_collect_read_timeout_sec or 120) cap = int(settings.ne_collect_run_timeout_cap_sec or 600) @@ -174,6 +176,7 @@ def _sample_one_target(target_row_id: str) -> None: row.last_error = "unsupported_vendor" db.commit() return + vendor_key = cmds.vendor_key cmd = detail_command(cmds, ifname) finally: db.close() @@ -190,8 +193,15 @@ def _sample_one_target(target_row_id: str) -> None: _set_target_error(target_row_id, _format_error(exc)) return - parsed = parse_zte_interface_detail(raw) - if parsed.in_bps == 0 and parsed.out_bps == 0 and parsed.bw_bps == 0 and not parsed.ifname: + parsed = parse_interface_detail(raw, vendor_key) + if ( + parsed.in_bps == 0 + and parsed.out_bps == 0 + and parsed.bw_bps == 0 + and parsed.in_util_pct == 0 + and parsed.out_util_pct == 0 + and not parsed.ifname + ): _set_target_error(target_row_id, "parse_empty") return diff --git a/netx_api/port_traffic_service.py b/netx_api/port_traffic_service.py index b991b7d..71b2f8d 100644 --- a/netx_api/port_traffic_service.py +++ b/netx_api/port_traffic_service.py @@ -15,8 +15,9 @@ from .cli_resolve import resolve_cli_target from .config import settings from .models import PortTrafficSample, PortTrafficTarget, PortTrafficTask from .ne_session_factory import close_netmiko_connection, open_netmiko_connection +from .ne_netmiko import send_show_command from .port_traffic_commands import commands_for_vendor -from .port_traffic_parsers import brief_port_to_dict, parse_zte_interface_brief +from .port_traffic_parsers import brief_port_to_dict, parse_interface_brief from .port_traffic_schemas import ( DiscoverPortItem, DiscoverPortsRequest, @@ -84,7 +85,7 @@ def _target_out(row: PortTrafficTarget) -> PortTrafficTargetOut: ) -def _assert_zte_targets(targets: list[PortTrafficTargetIn]) -> None: +def _assert_supported_targets(targets: list[PortTrafficTargetIn]) -> None: for t in targets: cmds = commands_for_vendor(t.vendor or "", "") if cmds is None: @@ -114,7 +115,7 @@ def get_task(db: Session, task_id: str) -> PortTrafficTaskOut: def create_task(db: Session, body: PortTrafficTaskCreate) -> PortTrafficTaskOut: - _assert_zte_targets(body.targets) + _assert_supported_targets(body.targets) now = _utcnow() status = "running" if body.start_now and body.targets else "draft" if body.start_now and not body.targets: @@ -219,7 +220,7 @@ def put_targets(db: Session, task_id: str, body: PortTrafficTargetsPut) -> list[ raise HTTPException(status_code=404, detail="task_not_found") if bool(task.collect_running): raise HTTPException(status_code=409, detail="collect_running") - _assert_zte_targets(body.targets) + _assert_supported_targets(body.targets) old = db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == task_id).all() old_ids = [str(t.id) for t in old] if old_ids: @@ -291,11 +292,14 @@ def discover_ports(db: Session, body: DiscoverPortsRequest) -> DiscoverPortsResp per_cmd = int(settings.ne_collect_read_timeout_sec or 120) conn = open_netmiko_connection(creds, session_timeout=per_cmd + 60) try: - raw = str(conn.send_command(command_string=cmds.brief, read_timeout=per_cmd) or "") + raw = send_show_command(conn, cmds.brief, read_timeout=per_cmd) finally: close_netmiko_connection(conn) - ports = [DiscoverPortItem(**brief_port_to_dict(p)) for p in parse_zte_interface_brief(raw)] + ports = [ + DiscoverPortItem(**brief_port_to_dict(p)) + for p in parse_interface_brief(raw, cmds.vendor_key) + ] return DiscoverPortsResponse( source=body.source, id=body.id, diff --git a/tests/test_port_traffic_parsers.py b/tests/test_port_traffic_parsers.py index ef4ddcb..391d26f 100644 --- a/tests/test_port_traffic_parsers.py +++ b/tests/test_port_traffic_parsers.py @@ -1,4 +1,4 @@ -"""Unit tests for ZTE port traffic brief/detail parsers (sample from oclaw Untitled-1.ps1).""" +"""Unit tests for ZTE / Huawei port traffic brief/detail parsers.""" from __future__ import annotations @@ -7,6 +7,12 @@ import unittest from netx_api.port_traffic_commands import commands_for_vendor, detail_command from netx_api.port_traffic_parsers import ( parse_bw_to_bps, + parse_cisco_interface_brief, + parse_cisco_interface_detail, + parse_huawei_interface_brief, + parse_huawei_interface_detail, + parse_interface_brief, + parse_interface_detail, parse_zte_interface_brief, parse_zte_interface_detail, ) @@ -56,6 +62,59 @@ smartgroup10 N/A N/A 1G up up up bvi2 N/A N/A up up up """ +HUAWEI_BRIEF = """\ +display interface brief +PHY: Physical +*down: administratively down +InUti/OutUti: input utility/output utility +Interface PHY Protocol InUti OutUti inErrors outErrors +Ethernet1/0/0 up up 0.01% 0% 0 0 +Ethernet1/0/1 up down 0% 0% 0 0 +GigabitEthernet0/0/0 up down 0% 0% 0 0 +NULL0 up up(s) 0% 0% 0 0 + +""" + +HUAWEI_DETAIL = """\ +display interface Ethernet1/0/0 +Ethernet1/0/0 current state : UP (ifindex: 5) +Line protocol current state : UP +Description: +Route Port,The Maximum Transmit Unit is 1500 + Last 300 seconds input rate: 0 bits/sec, 1 packets/sec + Last 300 seconds output rate: 0 bits/sec, 0 packets/sec + Input peak rate 0 bits/sec, Record time: - + Output peak rate 0 bits/sec, Record time: - + Last 300 seconds input utility rate: 0.01% + Last 300 seconds output utility rate: 0.00% + +""" + +CISCO_BRIEF = """\ +R2#show ip interface brief +Interface IP-Address OK? Method Status Protocol +GigabitEthernet0/0 192.168.0.128 YES NVRAM up up +GigabitEthernet0/1 172.16.0.2 YES manual up up +GigabitEthernet0/2 unassigned YES NVRAM administratively down down +GigabitEthernet0/3 unassigned YES NVRAM administratively down down +R2# +""" + +CISCO_DETAIL = """\ +R2#show interfaces GigabitEthernet0/1 +GigabitEthernet0/1 is up, line protocol is up + Hardware is iGbE, address is 5000.0003.0001 (bia 5000.0003.0001) + Internet address is 172.16.0.2/30 + MTU 1500 bytes, BW 1000000 Kbit/sec, DLY 10 usec, + reliability 255/255, txload 1/255, rxload 1/255 + Encapsulation ARPA, loopback not set + Last clearing of "show interface" counters never + 5 minute input rate 0 bits/sec, 0 packets/sec + 5 minute output rate 0 bits/sec, 0 packets/sec + 524 packets input, 158824 bytes, 0 no buffer +R2# +""" + class BwParseTests(unittest.TestCase): def test_compact(self): @@ -68,6 +127,9 @@ class BwParseTests(unittest.TestCase): def test_detail_line(self): self.assertEqual(parse_bw_to_bps("BW 1 Gbit/s"), 1_000_000_000) + def test_cisco_kbit_sec(self): + self.assertEqual(parse_bw_to_bps("MTU 1500 bytes, BW 1000000 Kbit/sec, DLY 10 usec,"), 1_000_000_000) + class BriefParserTests(unittest.TestCase): def test_sample_brief(self): @@ -96,7 +158,6 @@ class BriefParserTests(unittest.TestCase): self.assertTrue(any(n.startswith("xgei-") for n in names)) def test_rejects_prompt_without_hash(self): - # Some captures leave hostname alone after the table. text = ( SAMPLE_BRIEF + "\nAL5458-ACC-6120HS\n" @@ -119,17 +180,104 @@ class DetailParserTests(unittest.TestCase): self.assertEqual(d.in_util_pct, 0.0) self.assertEqual(d.out_util_pct, 0.0) self.assertEqual(d.description, "C2930L100-EQ2") - # Peak rates must not override Rate period values self.assertNotEqual(d.in_bps, 3536.0) +class HuaweiBriefParserTests(unittest.TestCase): + def test_sample_brief(self): + rows = parse_huawei_interface_brief(HUAWEI_BRIEF) + by_name = {r.ifname: r for r in rows} + self.assertEqual(len(rows), 4) + self.assertEqual(by_name["Ethernet1/0/0"].phy, "up") + self.assertEqual(by_name["Ethernet1/0/0"].prot, "up") + self.assertEqual(by_name["Ethernet1/0/0"].bw_bps, 0) + self.assertEqual(by_name["Ethernet1/0/1"].prot, "down") + self.assertEqual(by_name["NULL0"].prot, "up") # up(s) + self.assertEqual(by_name["GigabitEthernet0/0/0"].ifname, "GigabitEthernet0/0/0") + + def test_dispatch(self): + rows = parse_interface_brief(HUAWEI_BRIEF, "huawei") + self.assertTrue(any(r.ifname == "Ethernet1/0/0" for r in rows)) + + +class HuaweiDetailParserTests(unittest.TestCase): + def test_sample_detail(self): + d = parse_huawei_interface_detail(HUAWEI_DETAIL) + self.assertEqual(d.ifname, "Ethernet1/0/0") + self.assertEqual(d.admin_oper, "up") + self.assertEqual(d.bw_bps, 0) + self.assertEqual(d.rate_period_sec, 300) + self.assertEqual(d.in_bps, 0.0) + self.assertEqual(d.out_bps, 0.0) + self.assertEqual(d.in_util_pct, 0.01) + self.assertEqual(d.out_util_pct, 0.0) + + def test_dispatch(self): + d = parse_interface_detail(HUAWEI_DETAIL, "huawei") + self.assertEqual(d.in_util_pct, 0.01) + + +class CiscoBriefParserTests(unittest.TestCase): + def test_sample_brief(self): + rows = parse_cisco_interface_brief(CISCO_BRIEF) + by_name = {r.ifname: r for r in rows} + self.assertEqual(len(rows), 4) + self.assertEqual(by_name["GigabitEthernet0/0"].phy, "up") + self.assertEqual(by_name["GigabitEthernet0/0"].prot, "up") + self.assertEqual(by_name["GigabitEthernet0/1"].admin, "up") + self.assertEqual(by_name["GigabitEthernet0/2"].admin, "down") + self.assertEqual(by_name["GigabitEthernet0/2"].phy, "down") + self.assertEqual(by_name["GigabitEthernet0/2"].prot, "down") + + def test_dispatch(self): + rows = parse_interface_brief(CISCO_BRIEF, "cisco") + self.assertTrue(any(r.ifname == "GigabitEthernet0/1" for r in rows)) + + +class CiscoDetailParserTests(unittest.TestCase): + def test_sample_detail(self): + d = parse_cisco_interface_detail(CISCO_DETAIL) + self.assertEqual(d.ifname, "GigabitEthernet0/1") + self.assertEqual(d.admin_oper, "up") + self.assertEqual(d.bw_bps, 1_000_000_000) + self.assertEqual(d.rate_period_sec, 300) + self.assertEqual(d.in_bps, 0.0) + self.assertEqual(d.out_bps, 0.0) + self.assertEqual(d.in_util_pct, 0.0) + self.assertEqual(d.out_util_pct, 0.0) + + def test_dispatch(self): + d = parse_interface_detail(CISCO_DETAIL, "cisco") + self.assertEqual(d.bw_bps, 1_000_000_000) + + class CommandsTests(unittest.TestCase): def test_zte_matrix(self): cmds = commands_for_vendor("ZTE", "zxros") assert cmds is not None self.assertEqual(cmds.brief, "show interface brief") self.assertEqual(detail_command(cmds, "xgei-1/1/0/1"), "show interface xgei-1/1/0/1") - self.assertIsNone(commands_for_vendor("Cisco", "ios")) + self.assertIsNone(commands_for_vendor("Nokia", "sros")) + + def test_huawei_matrix(self): + cmds = commands_for_vendor("Huawei", "huawei_vrp") + assert cmds is not None + self.assertEqual(cmds.vendor_key, "huawei") + self.assertEqual(cmds.brief, "display interface brief") + self.assertEqual( + detail_command(cmds, "Ethernet1/0/0"), + "display interface Ethernet1/0/0", + ) + + def test_cisco_matrix(self): + cmds = commands_for_vendor("Cisco", "cisco_ios") + assert cmds is not None + self.assertEqual(cmds.vendor_key, "cisco") + self.assertEqual(cmds.brief, "show ip interface brief") + self.assertEqual( + detail_command(cmds, "GigabitEthernet0/1"), + "show interfaces GigabitEthernet0/1", + ) if __name__ == "__main__": diff --git a/web/WEB.md b/web/WEB.md index b69933a..f22ff78 100644 --- a/web/WEB.md +++ b/web/WEB.md @@ -108,7 +108,7 @@ src/ ## 端口流量监控 - API:`/v1/port-traffic/*`(任务 CRUD、discover/ports、samples、dashboard) -- 厂商:本期仅 ZTE(`show interface brief` / `show interface {if}`),解析 Rate period **bit/s** +- 厂商:ZTE(`show interface brief` / `show interface {if}`)、华为(`display interface brief` / `display interface {if}`)、思科(`show ip interface brief` / `show interfaces {if}`);解析速率 **bit/s**(华为暂无带宽则 `bw_bps=0`;思科暂无利用率则 util=0) - 调度:`NETX_PORT_TRAFFIC_SCHEDULER_ENABLED`(默认开),tick `NETX_PORT_TRAFFIC_SCHEDULER_TICK_SEC`(默认 15) - 单飞:同一任务同时只允许一轮采集;崩溃启动清除 `collect_running` - 保留:按任务 `retention_days` 清理过期 sample