From 31e745ef95359585c8339e027ba1bf91f2ce1327 Mon Sep 17 00:00:00 2001 From: oliver Date: Sat, 23 May 2026 11:19:00 +0800 Subject: [PATCH] fix(ume): show WSS connection state separately from alarm activity Track ws_connection for UI pills, stop overwriting runtime last_error on alarm events, and wire pause/resume to reconnect WSS. Co-authored-by: Cursor --- netx_api/main.py | 13 ++++- netx_api/ume_alarm_ws.py | 114 +++++++++++++++++++++++++++++--------- tests/test_ume_sync.py | 13 +++++ web/src/pages/UmePage.tsx | 34 ++++++++---- web/src/types.ts | 7 +++ 5 files changed, 142 insertions(+), 39 deletions(-) diff --git a/netx_api/main.py b/netx_api/main.py index 504bde0..00ec624 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -38,8 +38,10 @@ from .ume_alarm_ws import ( cancel_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, + get_ws_connection_status, get_ws_logs, load_persisted_subscription, + request_ws_reconnect, shutdown_ws_consumer, start_ume_alarm_ws_consumer, ) @@ -935,6 +937,7 @@ def ume_alarm_subscription_status(limit: int = 80) -> dict[str, Any]: return { "ok": True, **st, + "ws_connection": get_ws_connection_status(), "ws_consumer_status": str(ws_task.get("status") or ""), "ws_consumer_last_error": str(ws_task.get("last_error") or ""), "ws_consumer_last_run_at": ws_task.get("last_run_at"), @@ -1097,6 +1100,8 @@ def ume_runtime_task_pause(task: str) -> dict[str, Any]: _runtime_pause_task(tid) if tid in ("alarms_current_auto_sync", "inventory_auto_sync"): _clear_force_resume_hints(tid) + if tid == "alarms_current_ws_consumer": + request_ws_reconnect() _set_runtime_task(tid, status="paused", last_error="") return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()} @@ -1109,7 +1114,13 @@ def ume_runtime_task_resume(task: str) -> dict[str, Any]: _runtime_resume_task(tid) if tid in ("alarms_current_auto_sync", "inventory_auto_sync"): _request_force_sync_after_resume(tid) - _set_runtime_task(tid, status="running", last_error="已恢复:将跳过本轮周期等待并尽快同步") + resume_hint = "已恢复:将跳过本轮周期等待并尽快同步" + elif tid == "alarms_current_ws_consumer": + request_ws_reconnect() + resume_hint = "已恢复:将尽快重连 WSS" + else: + resume_hint = "已恢复" + _set_runtime_task(tid, status="running", last_error=resume_hint) return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()} diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py index 78d524c..a7177a4 100644 --- a/netx_api/ume_alarm_ws.py +++ b/netx_api/ume_alarm_ws.py @@ -90,11 +90,59 @@ def parse_subscription_id_from_already_exists_error(message: str) -> str: _shutdown_lock = threading.Lock() _subscription_lock = threading.Lock() _ws_wake_event = threading.Event() +_ws_connection_lock = threading.Lock() _subscription_id: str = "" _subscription_uri: str = "" _subscription_topic: str = "ALARM" _active_client: UMEClient | None = None +_ws_connection_state: str = "init" +_ws_connection_detail: str = "" + +_WS_CONNECTION_LABELS: dict[str, str] = { + "init": "初始化", + "connected": "已连接", + "connecting": "连接中", + "disconnected": "已断开", + "no_subscription": "无订阅", + "waiting_token": "等待 token", + "paused": "已暂停", + "error": "连接异常", + "reconnecting": "重连等待", +} + + +def _set_ws_connection_state(state: str, *, detail: str = "") -> None: + global _ws_connection_state, _ws_connection_detail + with _ws_connection_lock: + _ws_connection_state = str(state or "").strip() or "unknown" + _ws_connection_detail = str(detail or "").strip()[:240] + + +def get_ws_connection_status() -> dict[str, Any]: + with _ws_connection_lock: + state = str(_ws_connection_state or "init") + detail = str(_ws_connection_detail or "") + return { + "state": state, + "label": _WS_CONNECTION_LABELS.get(state, state), + "detail": detail, + } + + +def _notify_ws_connection( + state: str, + *, + detail: str = "", + subscription_id: str = "", + on_status: Callable[[str], None] | None = None, + log_level: str = "info", +) -> None: + _set_ws_connection_state(state, detail=detail) + label = _WS_CONNECTION_LABELS.get(state, state) + append_ws_log(detail or label, level=log_level, subscription_id=subscription_id) + if on_status is not None: + on_status(label) def _parse_ws_message(raw: str) -> dict[str, Any] | None: @@ -302,32 +350,34 @@ def _run_ws_session( label = _WS_ALARM_ACTION_LABEL.get(action, action) status_msg = f"{label} key={alarm_key}" + ("" if changed else " (无变更)") append_ws_log(status_msg, subscription_id=subscription_id) - if on_status is not None: - on_status(f"last={action}") except Exception as exc: db.rollback() append_ws_log(f"alarm apply failed: {str(exc)[:200]}", level="error", subscription_id=subscription_id) - if on_status is not None: - on_status(f"apply_error:{str(exc)[:120]}") finally: db.close() def _on_error(_ws: Any, error: Any) -> None: - append_ws_log(f"ws error: {str(error)[:200]}", level="error", subscription_id=subscription_id) - if on_status is not None: - on_status(f"ws_error:{str(error)[:120]}") + detail = f"ws error: {str(error)[:200]}" + _notify_ws_connection( + "error", + detail=detail, + subscription_id=subscription_id, + on_status=on_status, + log_level="error", + ) def _on_close(_ws: Any, close_status_code: Any, close_msg: Any) -> None: - append_ws_log( - f"ws closed code={close_status_code} msg={str(close_msg or '')[:120]}", - subscription_id=subscription_id, - ) + detail = f"ws closed code={close_status_code} msg={str(close_msg or '')[:120]}" + _notify_ws_connection("disconnected", detail=detail, subscription_id=subscription_id, on_status=on_status) closed.set() def _on_open(_ws: Any) -> None: - append_ws_log(f"ws connected uri={wss_uri[:120]}", subscription_id=subscription_id) - if on_status is not None: - on_status("connected") + _notify_ws_connection( + "connected", + detail=f"ws connected uri={wss_uri[:120]}", + subscription_id=subscription_id, + on_status=on_status, + ) ws_app = websocket.WebSocketApp( wss_uri, @@ -390,31 +440,41 @@ def run_alarm_ws_consumer_loop( backoff_s = 2.0 max_backoff_s = 120.0 - def _status(msg: str, *, level: str = "info", sub_id: str = "") -> None: - append_ws_log(msg, level=level, subscription_id=sub_id) - if on_status is not None: - on_status(msg) + def _loop_status( + state: str, + *, + detail: str = "", + subscription_id: str = "", + log_level: str = "info", + ) -> None: + _notify_ws_connection( + state, + detail=detail, + subscription_id=subscription_id, + on_status=on_status, + log_level=log_level, + ) while stop_event is None or not stop_event.is_set(): if is_paused is not None and is_paused(): - _status("paused", level="warning") + _loop_status("paused", log_level="warning") time.sleep(1.0) continue subscription_id, wss_uri = get_active_subscription() if not subscription_id or not wss_uri: - _status("no_subscription") + _loop_status("no_subscription") _ws_wake_event.wait(timeout=2.0) _ws_wake_event.clear() continue try: if not _wait_for_shared_token(client, timeout_s=120.0): - _status("waiting_token", level="warning", sub_id=subscription_id) + _loop_status("waiting_token", log_level="warning", subscription_id=subscription_id) time.sleep(2.0) continue - append_ws_log(f"connecting {wss_uri[:160]}", subscription_id=subscription_id) + _loop_status("connecting", detail=f"connecting {wss_uri[:160]}", subscription_id=subscription_id) _run_ws_session( client, wss_uri=wss_uri, @@ -425,21 +485,21 @@ def run_alarm_ws_consumer_loop( backoff_s = 2.0 except RuntimeError as exc: if "ume_ws_no_valid_token" in str(exc): - _status("waiting_token", level="warning", sub_id=subscription_id) + _loop_status("waiting_token", log_level="warning", subscription_id=subscription_id) time.sleep(2.0) continue - _status(f"session error: {str(exc)[:200]}", level="error", sub_id=subscription_id) + _loop_status("error", detail=f"session error: {str(exc)[:200]}", subscription_id=subscription_id, log_level="error") except Exception as exc: - _status(f"session error: {str(exc)[:200]}", level="error", sub_id=subscription_id) + _loop_status("error", detail=f"session error: {str(exc)[:200]}", subscription_id=subscription_id, log_level="error") sub_id_after, uri_after = get_active_subscription() if stop_event is not None and stop_event.is_set(): break if not sub_id_after or not uri_after: - _status("no_subscription") + _loop_status("no_subscription") continue - _status(f"reconnect in {int(backoff_s)}s", sub_id=sub_id_after) + _loop_status("reconnecting", detail=f"reconnect in {int(backoff_s)}s", subscription_id=sub_id_after) slept = 0.0 while slept < backoff_s: if stop_event is not None and stop_event.is_set(): diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 1464f95..303abcf 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -20,6 +20,7 @@ from netx_api.ume_alarm_ws import ( cancel_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, + get_ws_connection_status, get_ws_logs, load_persisted_subscription, parse_subscription_id_from_already_exists_error, @@ -738,6 +739,18 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertTrue(changed) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-1")) + def test_ws_connection_status_not_alarm_action(self): + from netx_api import ume_alarm_ws as ws_mod + + with ws_mod._ws_connection_lock: + ws_mod._ws_connection_state = "init" + ws_mod._ws_connection_detail = "" + ws_mod._set_ws_connection_state("connected", detail="ws connected") + st = get_ws_connection_status() + self.assertEqual(st["state"], "connected") + self.assertEqual(st["label"], "已连接") + self.assertEqual(st["detail"], "ws connected") + def test_ws_log_ring_buffer(self): from netx_api import ume_alarm_ws as ws_mod diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index da012c5..f4fcf21 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -197,6 +197,24 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { syncStatusQuery.data?.alarm_subscription ?? ({ active: false } as const); const wsConsumer = runtimeTasks.find((t) => t.task === "alarms_current_ws_consumer"); + const wsConn = subscriptionStatusQuery.data?.ws_connection; + const wsState = String(wsConn?.state || ""); + const wsPaused = Boolean(wsConsumer?.paused); + const wsLabel = + wsConn?.label || + (wsPaused ? "已暂停" : wsConsumer?.last_error || wsConsumer?.status || "-"); + const wsPillLevel = + wsPaused || wsState === "paused" + ? "warn" + : wsState === "connected" + ? "up" + : wsState === "connecting" || wsState === "reconnecting" || wsState === "waiting_token" + ? "warn" + : wsState === "no_subscription" || wsState === "disconnected" || wsState === "error" || wsState === "init" + ? "down" + : String(wsConsumer?.last_error || "").includes("connected") + ? "up" + : "unknown"; const wsLogs = [...(subscriptionStatusQuery.data?.ws_logs ?? alarmSub.ws_logs ?? [])].reverse(); const subscriptionActive = Boolean(alarmSub.active); const subPending = @@ -298,6 +316,7 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {

UME 告警订阅(WebSocket)

订阅需手动建立/取消;建立后(或重启后若库中仍有有效订阅)后台会自动连接 WSS 接收实时告警。 + 后台任务 alarms_current_ws_consumer 的「暂停/开始」会停止或恢复 WSS 连接(不会取消 UME 订阅)。

@@ -308,19 +327,12 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { id: {String(alarmSub.subscription_id).slice(0, 12)}… ) : null} - {wsConsumer ? ( + {subscriptionActive || wsConsumer ? ( - WSS: {wsConsumer.last_error || wsConsumer.status || "-"} + WSS: {wsLabel} ) : null}
diff --git a/web/src/types.ts b/web/src/types.ts index f25683c..f6f9cab 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -96,6 +96,12 @@ export type UmeWsLogEntry = { subscription_id?: string; }; +export type UmeWsConnectionStatus = { + state: string; + label: string; + detail?: string; +}; + export type UmeAlarmSubscriptionStatus = { ok?: boolean; created?: boolean; @@ -104,6 +110,7 @@ export type UmeAlarmSubscriptionStatus = { subscription_id?: string; wss_uri?: string; topic?: string; + ws_connection?: UmeWsConnectionStatus; ws_consumer_status?: string; ws_consumer_last_error?: string; ws_consumer_last_run_at?: string | null;