From 5fa73c9a4fa529f3843910810be4e71083e32bf7 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 30 Jul 2026 09:49:31 +0800 Subject: [PATCH] fix(webcrt): close CLI hop sessions when target exit returns to proxy Nested Huawei/ZTE/Cisco jumps keep the hop channel open after quit/exit; detect nested close or hop prompt return and tear down WebCRT so users cannot operate the proxy. Co-authored-by: Cursor --- netx_api/ne_session_factory.py | 125 ++++++++++++++++++++++++++++++++- netx_api/webcrt_service.py | 52 +++++++++++++- tests/test_cli_hop_return.py | 93 ++++++++++++++++++++++++ tests/test_webcrt.py | 75 ++++++++++++++++++++ 4 files changed, 343 insertions(+), 2 deletions(-) create mode 100644 tests/test_cli_hop_return.py diff --git a/netx_api/ne_session_factory.py b/netx_api/ne_session_factory.py index b990f70..97f7693 100644 --- a/netx_api/ne_session_factory.py +++ b/netx_api/ne_session_factory.py @@ -355,6 +355,101 @@ def _send_line(conn: ConnectHandler, line: str) -> None: conn.write_channel(text) +# Nested stelnet/telnet/ssh on vendor hops ends with messages like these; outer hop stays up. +_CLI_HOP_NESTED_END_RE = re.compile( + r"(?is)" + r"(?:^|\n)\s*(?:" + r"connection\s+closed(?:\s+by\s+(?:foreign|remote)\s+host)?" + r"|closed\s+by\s+foreign\s+host" + r"|connection\s+to\s+\S+\s+closed" + r"|%\s*connection\s+closed(?:\s+by\s+(?:foreign|remote)\s+host)?" + r"|\[connection\s+to\s+[^\]]+closed\]" + r"|remote\s+host\s+closed\s+the\s+connection" + r")[^\n]*\s*(?:\n|$)" +) + +# Last-line CLI prompts: [HW] Router# Router> +_CLI_PROMPT_LINE_RE = re.compile( + r"^(?:" + r"<[^>\r\n]{1,64}>|" + r"\[[^\]\r\n]{1,64}\]|" + r"[A-Za-z0-9][\w.\-:/]{0,62}[#>]" + r")\s*$" +) + + +def extract_cli_prompt_marker(text: str) -> str: + """Return the last recognizable CLI prompt line from channel text.""" + s = str(text or "").replace("\r\n", "\n").replace("\r", "\n") + for line in reversed(s.split("\n")): + # Strip common ANSI CSI sequences so markers match live reader bytes. + cleaned = re.sub(r"\x1b\[[0-9;?]*[A-Za-z]", "", line).strip() + if cleaned and _CLI_PROMPT_LINE_RE.match(cleaned): + return cleaned + return "" + + +def cli_hop_nested_session_ended(text: str) -> bool: + """True when nested jump (stelnet/telnet/ssh) reports connection closed.""" + return bool(_CLI_HOP_NESTED_END_RE.search(str(text or ""))) + + +def cli_hop_returned_to_proxy(text: str, hop_prompt: str) -> bool: + """True when output ends on the hop prompt captured before the jump command.""" + marker = str(hop_prompt or "").strip() + if not marker: + return False + last = extract_cli_prompt_marker(text) + return bool(last) and last == marker + + +def should_close_cli_hop_session( + recent: str, + hop_prompt: str = "", + *, + seen_other_prompt: bool = False, +) -> bool: + """Policy: end WebCRT when nested target session drops back to the hop CLI. + + Nested-close messages are matched only in a trailing window so a mid-session + ``display log`` that reprints old "Connection closed" text does not trip. + + Prompt-only return requires ``seen_other_prompt`` so identical default sysnames + (e.g. hop and target both ````) do not close immediately after jump. + """ + text = str(recent or "") + if cli_hop_nested_session_ended(text[-800:]): + return True + if not seen_other_prompt: + return False + return cli_hop_returned_to_proxy(text, hop_prompt) + + +def get_cli_hop_guard(conn: ConnectHandler | None) -> dict[str, Any] | None: + """Metadata attached by CLI hop connect; None when not a vendor CLI hop session.""" + if conn is None: + return None + guard = getattr(conn, "_netx_cli_hop", None) + if not isinstance(guard, dict) or not guard.get("enabled"): + return None + return guard + + +def _attach_cli_hop_guard( + conn: ConnectHandler, + *, + hop_prompt: str, + hop_vendor: str, + hop_host: str, +) -> None: + conn._netx_cli_hop = { # type: ignore[attr-defined] + "enabled": True, + "hop_prompt": str(hop_prompt or "").strip(), + "hop_vendor": str(hop_vendor or "").strip().lower(), + "hop_host": str(hop_host or "").strip(), + } + + 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)) @@ -429,10 +524,38 @@ def _connect_via_cli_hop( ) conn = ConnectHandler(**hop_dev) try: - _read_channel(conn, wait=0.5) + pre = _read_channel(conn, wait=0.5) + hop_prompt = extract_cli_prompt_marker(pre) + if not hop_prompt: + # Nudge hop CLI once so the prompt is visible for later return-to-proxy detection. + try: + conn.write_channel(getattr(conn, "RETURN", None) or "\n") + except Exception: + _send_line(conn, "") + pre = pre + _read_channel(conn, wait=0.35, max_loops=10) + hop_prompt = extract_cli_prompt_marker(pre) 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"])) + _attach_cli_hop_guard( + conn, + hop_prompt=hop_prompt, + hop_vendor=_hop_vendor(creds), + hop_host=hop_host, + ) + if hop_prompt: + _log.info( + "cli hop guard armed vendor=%s hop=%s prompt=%r", + _hop_vendor(creds), + hop_host, + hop_prompt, + ) + else: + _log.warning( + "cli hop guard armed without hop prompt vendor=%s hop=%s (nested-close only)", + _hop_vendor(creds), + hop_host, + ) return conn except Exception: try: diff --git a/netx_api/webcrt_service.py b/netx_api/webcrt_service.py index 8b498f4..357d025 100644 --- a/netx_api/webcrt_service.py +++ b/netx_api/webcrt_service.py @@ -21,7 +21,13 @@ from sqlalchemy.orm import Session from .config import settings from .ne_crypto import CredentialCryptoError -from .ne_session_factory import close_netmiko_connection, open_netmiko_connection +from .ne_session_factory import ( + close_netmiko_connection, + extract_cli_prompt_marker, + get_cli_hop_guard, + open_netmiko_connection, + should_close_cli_hop_session, +) _log = logging.getLogger("netx.webcrt") @@ -240,9 +246,14 @@ class WebcrtSession: # Only the newest attach_gen may consume out_queue / mark detach. attach_gen: int = 0 out_queue: queue.Queue[bytes | None] = field(default_factory=queue.Queue) + # Vendor CLI hop (Huawei/ZTE/Cisco): close when nested target session returns to hop. + cli_hop_guard: bool = False + cli_hop_prompt: str = "" _reader: threading.Thread | None = field(default=None, repr=False) _write_lock: threading.Lock = field(default_factory=threading.Lock, repr=False) _stdout_lock: threading.Lock = field(default_factory=threading.Lock, repr=False) + _hop_scan_buf: str = field(default="", repr=False) + _cli_hop_seen_other_prompt: bool = field(default=False, repr=False) def touch(self) -> None: self.last_activity = time.time() @@ -336,6 +347,7 @@ class WebcrtSession: self.out_queue.put(None) return channel = getattr(conn, "remote_conn", None) + hop_return = False try: while not self.closed: chunk = b"" @@ -383,9 +395,42 @@ class WebcrtSession: if chunk: self.touch() self.out_queue.put(chunk) + if self.cli_hop_guard and self._note_cli_hop_output(chunk): + hop_return = True + notice = ( + "\r\n*** WebCRT: 目标会话已结束,已断开代理连接 " + "(target session ended; closing hop proxy) ***\r\n" + ) + self.out_queue.put(notice.encode("utf-8", errors="replace")) + break finally: + if hop_return and not self.closed: + # Prefer registry close for audit + remove; fall back to local close. + try: + close_session(self.session_id, reason="cli_hop_return") + except Exception: + self.close("cli_hop_return") self.out_queue.put(None) + def _note_cli_hop_output(self, chunk: bytes) -> bool: + """Accumulate stdout and return True when nested CLI hop has returned to proxy.""" + try: + text = chunk.decode("utf-8", errors="replace") + except Exception: + text = str(chunk) + self._hop_scan_buf = (self._hop_scan_buf + text)[-12000:] + # Track a prompt that differs from the hop so same-sysname labs still need + # an explicit nested-close message before we tear down. + marker = str(self.cli_hop_prompt or "").strip() + last = extract_cli_prompt_marker(self._hop_scan_buf) + if last and (not marker or last != marker): + self._cli_hop_seen_other_prompt = True + return should_close_cli_hop_session( + self._hop_scan_buf, + self.cli_hop_prompt, + seen_other_prompt=self._cli_hop_seen_other_prompt, + ) + def close(self, reason: str = "closed") -> None: if self.closed: return @@ -570,6 +615,7 @@ def create_session( except Exception: pass + hop_guard = get_cli_hop_guard(conn) sess = WebcrtSession( session_id=session_id, ne_id=target_id, @@ -585,6 +631,8 @@ def create_session( bootstrap_output=str(bootstrap or "").encode("utf-8", errors="replace"), # Only nudge a live prompt when transcript has no recognizable prompt yet. needs_live_prompt=not _looks_like_cli_prompt(bootstrap), + cli_hop_guard=bool(hop_guard), + cli_hop_prompt=str((hop_guard or {}).get("hop_prompt") or ""), ) # Keep bootstrap for WS attach replay; do not rely solely on out_queue (StrictMode remount). sess.start_reader() @@ -602,6 +650,8 @@ def create_session( source=str(device.get("source") or ""), hop_enabled=bool(creds.get("hop_enabled")), hop_vendor=str(creds.get("hop_vendor") or "") if creds.get("hop_enabled") else "", + cli_hop_guard=bool(hop_guard), + cli_hop_prompt=str((hop_guard or {}).get("hop_prompt") or ""), client=client or "", active=active_session_count(), ) diff --git a/tests/test_cli_hop_return.py b/tests/test_cli_hop_return.py new file mode 100644 index 0000000..307f2d9 --- /dev/null +++ b/tests/test_cli_hop_return.py @@ -0,0 +1,93 @@ +"""Unit tests for vendor CLI hop return-to-proxy detection.""" + +from __future__ import annotations + +import unittest +from unittest.mock import MagicMock, patch + +from netx_api.ne_session_factory import ( + cli_hop_nested_session_ended, + cli_hop_returned_to_proxy, + extract_cli_prompt_marker, + get_cli_hop_guard, + should_close_cli_hop_session, +) + + +class CliHopReturnDetectionTests(unittest.TestCase): + def test_extract_huawei_and_cisco_prompts(self) -> None: + self.assertEqual(extract_cli_prompt_marker("banner\n\n"), "") + self.assertEqual(extract_cli_prompt_marker("[BJ-CORE]\n"), "[BJ-CORE]") + self.assertEqual(extract_cli_prompt_marker("R1#"), "R1#") + self.assertEqual(extract_cli_prompt_marker("R1>"), "R1>") + self.assertEqual(extract_cli_prompt_marker(""), "") + + def test_extract_strips_ansi(self) -> None: + self.assertEqual( + extract_cli_prompt_marker("\x1b[32m\x1b[0m"), + "", + ) + + def test_nested_session_end_messages(self) -> None: + self.assertTrue(cli_hop_nested_session_ended("quit\nConnection closed by foreign host\n")) + self.assertTrue(cli_hop_nested_session_ended("Connection closed\n")) + self.assertTrue(cli_hop_nested_session_ended("% Connection closed by remote host\n")) + self.assertFalse(cli_hop_nested_session_ended("\n")) + + def test_returned_to_proxy_requires_exact_last_prompt(self) -> None: + self.assertTrue(cli_hop_returned_to_proxy("x\n\n", "")) + self.assertFalse(cli_hop_returned_to_proxy("x\n\n", "")) + self.assertFalse(cli_hop_returned_to_proxy("mentions in text\n\n", "")) + + def test_should_close_on_nested_end_in_tail(self) -> None: + old = ("old Connection closed by foreign host\n" * 40) + "\n" + self.assertFalse(should_close_cli_hop_session(old, "", seen_other_prompt=True)) + fresh = old + "Connection closed by foreign host\n\n" + self.assertTrue(should_close_cli_hop_session(fresh, "", seen_other_prompt=True)) + + def test_prompt_only_needs_seen_other_prompt(self) -> None: + text = "work\n\n" + self.assertFalse(should_close_cli_hop_session(text, "", seen_other_prompt=False)) + self.assertTrue(should_close_cli_hop_session(text, "", seen_other_prompt=True)) + + @patch("netx_api.ne_session_factory._interactive_target_auth") + @patch("netx_api.ne_session_factory._read_channel") + @patch("netx_api.ne_session_factory.ConnectHandler") + def test_connect_attaches_cli_hop_guard( + self, + mock_ch: MagicMock, + mock_read: MagicMock, + mock_auth: MagicMock, + ) -> None: + from netx_api.ne_session_factory import _connect_via_cli_hop + + conn = MagicMock() + mock_ch.return_value = conn + mock_read.side_effect = ["\n", ""] + mock_auth.return_value = None + creds = { + "hop_host": "10.0.0.1", + "hop_username": "admin", + "hop_password": "hop-pass", + "hop_protocol": "ssh", + "hop_vendor": "huawei", + "hop_port": 22, + "hop_vrf": "", + "hop_command_template": "", + "username": "target", + "password": "target-pass", + "ip_address": "10.0.0.2", + "port": 22, + } + out = _connect_via_cli_hop(creds) + self.assertIs(out, conn) + guard = get_cli_hop_guard(conn) + self.assertIsNotNone(guard) + assert guard is not None + self.assertTrue(guard["enabled"]) + self.assertEqual(guard["hop_prompt"], "") + self.assertEqual(guard["hop_vendor"], "huawei") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_webcrt.py b/tests/test_webcrt.py index 1cf151e..9ceb980 100644 --- a/tests/test_webcrt.py +++ b/tests/test_webcrt.py @@ -343,6 +343,81 @@ class WebcrtServiceTests(unittest.TestCase): svc._reap_sessions() self.assertIsNone(svc.get_session("idle")) + @patch.object(svc, "_audit") + def test_cli_hop_return_closes_session(self, _mock_audit: MagicMock) -> None: + """Vendor CLI hop: nested target exit must tear down WebCRT (no hop shell).""" + conn = _FakeConn() + chunks = [ + b"\r\n", + b"quit\r\nConnection closed by foreign host\r\n\r\n\r\n", + ] + idx = {"i": 0} + + def recv_ready() -> bool: + return idx["i"] < len(chunks) + + def recv(_n: int) -> bytes: + i = idx["i"] + idx["i"] = i + 1 + return chunks[i] + + conn.remote_conn.recv_ready.side_effect = recv_ready + conn.remote_conn.recv.side_effect = recv + conn.remote_conn.exit_status_ready.return_value = False + + sess = svc.WebcrtSession( + session_id="hop1", + ne_id="ne1", + ne_name="lab", + ne_ip="1.2.3.4", + protocol="ssh", + cols=80, + rows=24, + conn=conn, # type: ignore[arg-type] + cli_hop_guard=True, + cli_hop_prompt="", + ) + with svc._sessions_lock: + svc._sessions["hop1"] = sess + sess.start_reader() + deadline = time.time() + 3.0 + got: list[bytes] = [] + while time.time() < deadline: + item = sess.take_stdout(0, timeout=0.2) + if item == "empty": + if sess.closed: + break + continue + if item is None: + break + if isinstance(item, bytes): + got.append(item) + self.assertTrue(sess.closed) + self.assertEqual(sess.close_reason, "cli_hop_return") + self.assertIsNone(svc.get_session("hop1")) + blob = b"".join(got).decode("utf-8", errors="replace") + self.assertIn("Connection closed by foreign host", blob) + self.assertIn("目标会话已结束", blob) + + def test_cli_hop_note_ignores_same_sysname_until_close_msg(self) -> None: + sess = svc.WebcrtSession( + session_id="hop2", + ne_id="ne1", + ne_name="lab", + ne_ip="1.2.3.4", + protocol="ssh", + cols=80, + rows=24, + cli_hop_guard=True, + cli_hop_prompt="", + ) + # Same default sysname on target must not trip prompt-only close. + self.assertFalse(sess._note_cli_hop_output(b"\r\n")) + self.assertFalse(sess._cli_hop_seen_other_prompt) + self.assertTrue( + sess._note_cli_hop_output(b"Connection closed by foreign host\r\n\r\n") + ) + if __name__ == "__main__": unittest.main()