feat(ne): ZTE jump host for connect test and collection

Add hop fields and encrypted credentials, unified Netmiko session factory with ZTE ssh/telnet CLI templates and secondary auth, wire connect/collect paths, and NE management UI with auto-suggested jump commands.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-28 17:14:49 +08:00
parent 71ec3385d2
commit 4575fef523
13 changed files with 635 additions and 56 deletions

View file

@ -740,6 +740,15 @@ def on_startup() -> None:
conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS user_label")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS ne_name")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS user_label")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_enabled BOOLEAN DEFAULT FALSE")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_vendor VARCHAR(32) DEFAULT 'zte'")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_host VARCHAR(128) DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_port INTEGER DEFAULT 22")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_protocol VARCHAR(16) DEFAULT 'ssh'")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_username VARCHAR(128) DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_password_enc TEXT DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_command_template TEXT DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE managed_ne ADD COLUMN IF NOT EXISTS hop_vrf VARCHAR(128) DEFAULT ''")
conn.exec_driver_sql(
"ALTER TABLE ne_collection_job ADD COLUMN IF NOT EXISTS last_run_at TIMESTAMP"
)

View file

@ -244,6 +244,15 @@ class ManagedNE(Base):
site: Mapped[str] = mapped_column(String(256), default="")
tags: Mapped[str] = mapped_column(String(512), default="")
remark: Mapped[str] = mapped_column(String(1024), default="")
hop_enabled: Mapped[bool] = mapped_column(default=False)
hop_vendor: Mapped[str] = mapped_column(String(32), default="zte")
hop_host: Mapped[str] = mapped_column(String(128), default="")
hop_port: Mapped[int] = mapped_column(Integer, default=22)
hop_protocol: Mapped[str] = mapped_column(String(16), default="ssh")
hop_username: Mapped[str] = mapped_column(String(128), default="")
hop_password_enc: Mapped[str] = mapped_column(Text, default="")
hop_command_template: Mapped[str] = mapped_column(Text, default="")
hop_vrf: Mapped[str] = mapped_column(String(128), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)

View file

@ -9,16 +9,14 @@ from datetime import datetime
from pathlib import Path
from typing import Any
from netmiko import ConnectHandler
from .collection_job_state import finalize_collection_job, sync_job_progress
from .config import settings
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 normalize_netmiko_device_type
from .ne_service import get_device_credentials
from .ne_session_factory import open_netmiko_connection
_log = logging.getLogger("netx.ne.collect")
_executor: ThreadPoolExecutor | None = None
@ -45,32 +43,24 @@ def _safe_filename_part(text: str) -> str:
def _collect_on_device(creds: dict[str, Any], commands: list[str]) -> str:
device_type = normalize_netmiko_device_type(creds["device_type"], creds["protocol"])
per_cmd = int(settings.ne_collect_read_timeout_sec or 120)
dev: dict[str, Any] = {
"device_type": device_type,
"host": creds["ip_address"],
"username": creds["username"],
"password": creds["password"],
"port": int(creds["port"] or 22),
"conn_timeout": int(settings.ne_connect_timeout_sec or 30),
"auth_timeout": int(settings.ne_connect_timeout_sec or 30),
"banner_timeout": int(settings.ne_connect_timeout_sec or 30),
"session_timeout": per_cmd * max(1, len(commands)) + 60,
}
secret = str(creds.get("enable_secret") or "").strip()
if secret:
dev["secret"] = secret
chunks: list[str] = []
with ConnectHandler(**dev) as conn:
session_timeout = per_cmd * max(1, len(commands)) + 60
conn = open_netmiko_connection(creds, session_timeout=session_timeout)
try:
prompt = str(conn.find_prompt() or "")
chunks: list[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)
chunks.append(str(out or ""))
chunks.append("\n")
return "".join(chunks)
return "".join(chunks)
finally:
try:
conn.disconnect()
except Exception:
pass
def _collect_with_timeout(creds: dict[str, Any], commands: list[str]) -> str:

View file

@ -6,14 +6,12 @@ from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from typing import Any
from netmiko import ConnectHandler
from .config import settings
from .db import SessionLocal
from .models import ManagedNE
from .ne_crypto import CredentialCryptoError, decrypt_secret
from .ne_netmiko import normalize_netmiko_device_type
from .ne_crypto import CredentialCryptoError
from .ne_service import get_device_credentials
from .ne_session_factory import open_netmiko_connection
_log = logging.getLogger("netx.ne.connect")
_executor: ThreadPoolExecutor | None = None
@ -102,42 +100,53 @@ def _clean_prompt_hostname(prompt: str) -> str | None:
return p[:256]
def _classify_connect_error(creds: dict[str, Any], exc: BaseException) -> str:
raw = str(exc).lower()
detail = str(exc).split("\n")[0][:480]
if creds.get("hop_enabled"):
if "hop_credentials_incomplete" in raw or "hop_command_template_invalid" in raw:
return detail
if "target_auth_timeout" in raw:
return "target_auth_failed: " + detail
if "authentication" in raw or "auth" in raw:
if "hop_host" in raw or str(creds.get("hop_host") or "") in raw:
return "hop_auth_failed: " + detail
return "target_auth_failed: " + detail
if "timed out" in raw or "timeout" in raw:
return "hop_connect_failed: " + detail
return "hop_command_failed: " + detail
return detail
def _probe_device(creds: dict[str, Any]) -> tuple[str, str, str | None]:
"""Login via Netmiko, probe hostname, return (status, message, discovered_name)."""
device_type = normalize_netmiko_device_type(creds["device_type"], creds["protocol"])
vendor = str(creds.get("vendor") or "")
dev: dict[str, Any] = {
"device_type": device_type,
"host": creds["ip_address"],
"username": creds["username"],
"password": creds["password"],
"port": int(creds["port"] or 22),
"conn_timeout": int(settings.ne_connect_timeout_sec or 30),
"auth_timeout": int(settings.ne_connect_timeout_sec or 30),
"banner_timeout": int(settings.ne_connect_timeout_sec or 30),
}
secret = str(creds.get("enable_secret") or "").strip()
if secret:
dev["secret"] = secret
session_timeout = 180 if creds.get("hop_enabled") else None
conn = None
try:
with ConnectHandler(**dev) as conn:
prompt = str(conn.find_prompt() or "")
command = hostname_probe_command(creds["device_type"], vendor)
output = ""
if command:
output = conn.send_command(command_string=command, read_timeout=30)
hostname = parse_hostname_from_output(creds["device_type"], vendor, output, prompt)
if hostname:
return "pass", f"connected: {hostname}", hostname
if command:
return "pass", "connected (hostname not parsed)", None
fallback = _clean_prompt_hostname(prompt)
if fallback:
return "pass", f"connected: {fallback}", fallback
return "pass", "connected", None
conn = open_netmiko_connection(creds, session_timeout=session_timeout)
prompt = str(conn.find_prompt() or "")
command = hostname_probe_command(creds["device_type"], vendor)
output = ""
if command:
output = conn.send_command(command_string=command, read_timeout=30)
hostname = parse_hostname_from_output(creds["device_type"], vendor, output, prompt)
if hostname:
return "pass", f"connected: {hostname}", hostname
if command:
return "pass", "connected (hostname not parsed)", None
fallback = _clean_prompt_hostname(prompt)
if fallback:
return "pass", f"connected: {fallback}", fallback
return "pass", "connected", None
except Exception as exc:
msg = str(exc).split("\n")[0][:480]
return "fail", msg, None
return "fail", _classify_connect_error(creds, exc), None
finally:
if conn is not None:
try:
conn.disconnect()
except Exception:
pass
def _update_row(ne_id: str, status: str, message: str, discovered_name: str | None = None) -> None:

View file

@ -21,6 +21,15 @@ class ManagedNeCreate(BaseModel):
password: str
tags: str = ""
remark: str = ""
hop_enabled: bool = False
hop_vendor: str = "zte"
hop_host: str = ""
hop_port: int = 22
hop_protocol: str = "ssh"
hop_username: str = ""
hop_password: str = ""
hop_command_template: str = ""
hop_vrf: str = ""
@field_validator("vendor")
@classmethod
@ -45,6 +54,15 @@ class ManagedNeUpdate(BaseModel):
password: str | None = None
tags: str | None = None
remark: str | None = None
hop_enabled: bool | None = None
hop_vendor: str | None = None
hop_host: str | None = None
hop_port: int | None = None
hop_protocol: str | None = None
hop_username: str | None = None
hop_password: str | None = None
hop_command_template: str | None = None
hop_vrf: str | None = None
@field_validator("vendor")
@classmethod
@ -74,6 +92,14 @@ class ManagedNeOut(BaseModel):
connect_tested_at: datetime | None
tags: str
remark: str
hop_enabled: bool = False
hop_vendor: str = "zte"
hop_host: str = ""
hop_port: int = 22
hop_protocol: str = "ssh"
hop_username: str = ""
hop_command_template: str = ""
hop_vrf: str = ""
created_at: datetime
updated_at: datetime

View file

@ -43,6 +43,63 @@ def _normalize_protocol(protocol: str) -> str:
return p if p in ("ssh", "telnet") else "ssh"
def _normalize_hop_vendor(vendor: str) -> str:
v = str(vendor or "zte").strip().lower()
return v if v in ("zte",) else "zte"
def _validate_hop_on_create(body: ManagedNeCreate) -> None:
if not body.hop_enabled:
return
if not str(body.hop_host or "").strip():
raise HTTPException(status_code=400, detail="hop_host_required")
if not str(body.hop_username or "").strip():
raise HTTPException(status_code=400, detail="hop_username_required")
if not str(body.hop_password or "").strip():
raise HTTPException(status_code=400, detail="hop_password_required")
def _apply_hop_create(row: ManagedNE, body: ManagedNeCreate) -> None:
row.hop_enabled = bool(body.hop_enabled)
row.hop_vendor = _normalize_hop_vendor(body.hop_vendor)
row.hop_host = str(body.hop_host or "").strip()
row.hop_port = int(body.hop_port or 22)
row.hop_protocol = _normalize_protocol(body.hop_protocol)
row.hop_username = str(body.hop_username or "").strip()
row.hop_password_enc = encrypt_secret(body.hop_password) if body.hop_enabled else ""
row.hop_command_template = str(body.hop_command_template or "").strip()
row.hop_vrf = str(body.hop_vrf or "").strip()
def _apply_hop_update(row: ManagedNE, data: dict[str, Any]) -> None:
if "hop_enabled" in data and data["hop_enabled"] is not None:
row.hop_enabled = bool(data["hop_enabled"])
if "hop_vendor" in data and data["hop_vendor"] is not None:
row.hop_vendor = _normalize_hop_vendor(data["hop_vendor"])
if "hop_host" in data and data["hop_host"] is not None:
row.hop_host = str(data["hop_host"]).strip()
if "hop_port" in data and data["hop_port"] is not None:
row.hop_port = int(data["hop_port"])
if "hop_protocol" in data and data["hop_protocol"] is not None:
row.hop_protocol = _normalize_protocol(data["hop_protocol"])
if "hop_username" in data and data["hop_username"] is not None:
row.hop_username = str(data["hop_username"]).strip()
if "hop_password" in data and data["hop_password"]:
_require_crypto()
row.hop_password_enc = encrypt_secret(str(data["hop_password"]))
if "hop_command_template" in data and data["hop_command_template"] is not None:
row.hop_command_template = str(data["hop_command_template"]).strip()
if "hop_vrf" in data and data["hop_vrf"] is not None:
row.hop_vrf = str(data["hop_vrf"]).strip()
if row.hop_enabled:
if not str(row.hop_host or "").strip():
raise HTTPException(status_code=400, detail="hop_host_required")
if not str(row.hop_username or "").strip():
raise HTTPException(status_code=400, detail="hop_username_required")
if not str(row.hop_password_enc or "").strip():
raise HTTPException(status_code=400, detail="hop_password_required")
def row_to_out(row: ManagedNE) -> ManagedNeOut:
status = str(row.connect_status or "unknown")
if status not in ("unknown", "testing", "pass", "fail"):
@ -61,6 +118,14 @@ def row_to_out(row: ManagedNE) -> ManagedNeOut:
connect_tested_at=row.connect_tested_at,
tags=str(row.tags or ""),
remark=str(row.remark or ""),
hop_enabled=bool(row.hop_enabled),
hop_vendor=str(row.hop_vendor or "zte"),
hop_host=str(row.hop_host or ""),
hop_port=int(row.hop_port or 22),
hop_protocol=str(row.hop_protocol or "ssh"),
hop_username=str(row.hop_username or ""),
hop_command_template=str(row.hop_command_template or ""),
hop_vrf=str(row.hop_vrf or ""),
created_at=row.created_at,
updated_at=row.updated_at,
)
@ -114,6 +179,7 @@ def get_managed_ne(db: Session, ne_id: str) -> ManagedNeOut:
def create_managed_ne(db: Session, body: ManagedNeCreate) -> ManagedNeOut:
_require_crypto()
_validate_hop_on_create(body)
ip = _normalize_ip(body.ip_address)
if not ip:
raise HTTPException(status_code=400, detail="ip_address_required")
@ -139,6 +205,7 @@ def create_managed_ne(db: Session, body: ManagedNeCreate) -> ManagedNeOut:
created_at=now,
updated_at=now,
)
_apply_hop_create(row, body)
db.add(row)
db.commit()
db.refresh(row)
@ -180,6 +247,19 @@ def update_managed_ne(db: Session, ne_id: str, body: ManagedNeUpdate) -> Managed
if "password" in data and data["password"]:
_require_crypto()
row.password_enc = encrypt_secret(str(data["password"]))
hop_keys = (
"hop_enabled",
"hop_vendor",
"hop_host",
"hop_port",
"hop_protocol",
"hop_username",
"hop_password",
"hop_command_template",
"hop_vrf",
)
if any(k in data for k in hop_keys):
_apply_hop_update(row, data)
row.updated_at = _now()
db.commit()
db.refresh(row)
@ -311,6 +391,10 @@ def import_managed_ne(db: Session, content: bytes, filename: str) -> ImportResul
def get_device_credentials(row: ManagedNE) -> dict[str, Any]:
hop_enabled = bool(row.hop_enabled)
hop_password = ""
if hop_enabled and str(row.hop_password_enc or "").strip():
hop_password = decrypt_secret(row.hop_password_enc)
return {
"id": str(row.id),
"vendor": str(row.vendor or ""),
@ -322,4 +406,13 @@ def get_device_credentials(row: ManagedNE) -> dict[str, Any]:
"password": decrypt_secret(row.password_enc),
"enable_secret": decrypt_secret(row.enable_secret_enc),
"name": str(row.name or ""),
"hop_enabled": hop_enabled,
"hop_vendor": str(row.hop_vendor or "zte"),
"hop_host": str(row.hop_host or ""),
"hop_port": int(row.hop_port or 22),
"hop_protocol": str(row.hop_protocol or "ssh"),
"hop_username": str(row.hop_username or ""),
"hop_password": hop_password,
"hop_command_template": str(row.hop_command_template or ""),
"hop_vrf": str(row.hop_vrf or ""),
}

View file

@ -0,0 +1,189 @@
"""Netmiko session factory: direct connect or via ZTE jump host."""
from __future__ import annotations
import logging
import re
import time
from typing import Any
from netmiko import ConnectHandler
from .config import settings
from .ne_netmiko import normalize_netmiko_device_type
_log = logging.getLogger("netx.ne.session")
_HOP_PLACEHOLDERS = ("target_ip", "target_port", "target_user", "target_password", "vrf")
# ZTE CLI jump: ssh/telnet <ip> [vrf <name>] — target user/password via secondary auth.
_LEGACY_HOP_TEMPLATES = frozenset({"ssh {target_user}@{target_ip}", "ssh {target_ip}", "telnet {target_ip}"})
def default_zte_hop_template(protocol: str, vrf: str = "") -> str:
cmd = "telnet" if str(protocol or "ssh").strip().lower() == "telnet" else "ssh"
v = str(vrf or "").strip()
if v:
return f"{cmd} {{target_ip}} vrf {{vrf}}"
return f"{cmd} {{target_ip}}"
def render_hop_command(template: str, creds: dict[str, Any]) -> str:
"""Render hop command from template using whitelisted placeholders only."""
tpl = str(template or "").strip()
if not tpl or tpl in _LEGACY_HOP_TEMPLATES:
tpl = default_zte_hop_template(
str(creds.get("hop_protocol") or "ssh"),
str(creds.get("hop_vrf") or ""),
)
values = {
"target_ip": str(creds.get("ip_address") or ""),
"target_port": str(int(creds.get("port") or 22)),
"target_user": str(creds.get("username") or ""),
"target_password": str(creds.get("password") or ""),
"vrf": str(creds.get("hop_vrf") or "").strip(),
}
out = tpl
for key in _HOP_PLACEHOLDERS:
out = out.replace("{" + key + "}", values[key])
if "{" in out or "}" in out:
raise ValueError("hop_command_template_invalid_placeholder")
return out
def _base_connect_kwargs(
*,
device_type: str,
host: str,
port: int,
username: str,
password: str,
enable_secret: str,
session_timeout: int | None = None,
) -> dict[str, Any]:
timeout = int(settings.ne_connect_timeout_sec or 30)
dev: dict[str, Any] = {
"device_type": device_type,
"host": host,
"username": username,
"password": password,
"port": int(port or 22),
"conn_timeout": timeout,
"auth_timeout": timeout,
"banner_timeout": timeout,
}
if session_timeout is not None:
dev["session_timeout"] = session_timeout
secret = str(enable_secret or "").strip()
if secret:
dev["secret"] = secret
return dev
def _connect_direct(creds: dict[str, Any], *, session_timeout: int | None = None) -> ConnectHandler:
device_type = normalize_netmiko_device_type(creds["device_type"], creds["protocol"])
dev = _base_connect_kwargs(
device_type=device_type,
host=str(creds["ip_address"]),
port=int(creds["port"] or 22),
username=str(creds["username"]),
password=str(creds["password"]),
enable_secret=str(creds.get("enable_secret") or ""),
session_timeout=session_timeout,
)
return ConnectHandler(**dev)
def _read_channel(conn: ConnectHandler, wait: float = 0.5, max_loops: int = 40) -> str:
time.sleep(wait)
chunks: list[str] = []
for _ in range(max_loops):
part = conn.read_channel()
if not part:
break
chunks.append(part)
time.sleep(0.2)
return "".join(chunks)
def _send_line(conn: ConnectHandler, line: str) -> None:
text = str(line or "")
if not text.endswith("\n"):
text += "\n"
conn.write_channel(text)
def _prompt_needs_auth(text: str) -> tuple[bool, bool]:
low = text.lower()
need_user = bool(re.search(r"(username|login|user\s*name)\s*[:>]", low))
need_pass = bool(re.search(r"password\s*[:>]", low))
return need_user, need_pass
def _interactive_target_auth(conn: ConnectHandler, username: str, password: str) -> None:
"""Respond to username/password prompts after hop command (target credentials)."""
deadline = time.time() + int(settings.ne_connect_timeout_sec or 30)
sent_user = False
sent_pass = False
while time.time() < deadline:
buf = _read_channel(conn, wait=0.3, max_loops=8)
need_user, need_pass = _prompt_needs_auth(buf)
if need_pass and not sent_pass:
_send_line(conn, password)
sent_pass = True
continue
if need_user and not sent_user:
_send_line(conn, username)
sent_user = True
continue
if sent_pass and not need_user and not need_pass:
return
if not buf.strip():
time.sleep(0.3)
continue
if re.search(r"[>#]\s*$", buf):
if sent_pass or (sent_user and not need_pass):
return
time.sleep(0.3)
if not sent_pass:
raise TimeoutError("target_auth_timeout")
def _connect_via_zte_hop(creds: dict[str, Any], *, session_timeout: int | None = None) -> ConnectHandler:
hop_host = str(creds.get("hop_host") or "").strip()
hop_user = str(creds.get("hop_username") or "").strip()
hop_pass = str(creds.get("hop_password") or "")
if not hop_host or not hop_user or not hop_pass:
raise ValueError("hop_credentials_incomplete")
hop_protocol = str(creds.get("hop_protocol") or "ssh")
hop_device_type = normalize_netmiko_device_type("zte_zxros", hop_protocol)
hop_dev = _base_connect_kwargs(
device_type=hop_device_type,
host=hop_host,
port=int(creds.get("hop_port") or 22),
username=hop_user,
password=hop_pass,
enable_secret="",
session_timeout=session_timeout or 180,
)
conn = ConnectHandler(**hop_dev)
try:
_read_channel(conn, wait=0.5)
hop_cmd = render_hop_command(str(creds.get("hop_command_template") or ""), creds)
_send_line(conn, hop_cmd)
_interactive_target_auth(conn, str(creds["username"]), str(creds["password"]))
return conn
except Exception:
try:
conn.disconnect()
except Exception:
pass
raise
def open_netmiko_connection(creds: dict[str, Any], *, session_timeout: int | None = None) -> ConnectHandler:
"""Open a Netmiko connection to the target NE (direct or via configured hop)."""
if creds.get("hop_enabled"):
return _connect_via_zte_hop(creds, session_timeout=session_timeout)
return _connect_direct(creds, session_timeout=session_timeout)