From 5bd83c58dcb511647be99796217b72460d038eb6 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 22 May 2026 21:15:15 +0800 Subject: [PATCH] feat(ume): expose WSS runtime logs on subscription status UI Ring-buffer ws_logs API, alarm raise/clear labels with alarmkey, and scrollable log panel on UME page. Co-authored-by: Cursor --- netx_api/main.py | 5 +- netx_api/ume_alarm_ws.py | 131 +++++++++++++++++++++++++++++--------- tests/test_ume_sync.py | 18 ++++++ web/src/index.css | 40 ++++++++++++ web/src/pages/UmePage.tsx | 32 ++++++++++ web/src/types.ts | 8 +++ 6 files changed, 203 insertions(+), 31 deletions(-) diff --git a/netx_api/main.py b/netx_api/main.py index c819d56..504bde0 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -38,6 +38,7 @@ from .ume_alarm_ws import ( cancel_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, + get_ws_logs, load_persisted_subscription, shutdown_ws_consumer, start_ume_alarm_ws_consumer, @@ -927,15 +928,17 @@ def ume_token_disconnect() -> dict[str, Any]: @app.get("/v1/ume/alarm-subscription/status") -def ume_alarm_subscription_status() -> dict[str, Any]: +def ume_alarm_subscription_status(limit: int = 80) -> dict[str, Any]: st = get_subscription_status() ws_task = _UME_RUNTIME_TASKS.get("alarms_current_ws_consumer") or {} + log_limit = max(10, min(int(limit or 80), 100)) return { "ok": True, **st, "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"), + "ws_logs": get_ws_logs(limit=log_limit), } diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py index 7f3476a..78d524c 100644 --- a/netx_api/ume_alarm_ws.py +++ b/netx_api/ume_alarm_ws.py @@ -1,5 +1,7 @@ from __future__ import annotations +from collections import deque +from datetime import datetime, timezone import json import logging import re @@ -19,13 +21,67 @@ from .ume_alarm_subscription_store import ( save_subscription, ) from .ume_client import UMEClient -from .ume_sync_service import apply_alarm_to_current, extract_alarm_from_notification, _utc_now_naive +from .ume_sync_service import ( + _alarm_key, + apply_alarm_to_current, + extract_alarm_from_notification, + normalize_yang_alarm, + _utc_now_naive, +) + +_WS_ALARM_ACTION_LABEL: dict[str, str] = { + "inserted": "上报(新增)", + "updated": "上报(更新)", + "deleted": "清除", + "skipped": "忽略", +} _ws_log = logging.getLogger("netx.ume.alarm_ws") +_WS_LOG_LOCK = threading.Lock() +_WS_LOG_ENTRIES: deque[dict[str, Any]] = deque(maxlen=200) +_WS_LOG_MAX_RETURN = 100 + _ORPHAN_SUB_ID_RE = re.compile(r"id:([0-9a-fA-F-]{8}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{12})") +def _utc_now_iso() -> str: + return datetime.now(timezone.utc).isoformat() + + +def append_ws_log( + message: str, + *, + level: str = "info", + subscription_id: str = "", +) -> None: + """Ring buffer of recent WSS events for UI (newest last).""" + entry = { + "ts": _utc_now_iso(), + "level": str(level or "info").strip().lower() or "info", + "message": str(message or "").strip()[:500], + "subscription_id": str(subscription_id or "").strip(), + } + with _WS_LOG_LOCK: + _WS_LOG_ENTRIES.append(entry) + log_fn = _ws_log.info + if entry["level"] == "warning": + log_fn = _ws_log.warning + elif entry["level"] == "error": + log_fn = _ws_log.error + log_fn("%s%s", entry["message"], f" sub={entry['subscription_id']}" if entry["subscription_id"] else "") + + +def get_ws_logs(*, limit: int | None = None) -> list[dict[str, Any]]: + cap = int(limit if limit is not None else _WS_LOG_MAX_RETURN) + cap = max(1, min(cap, _WS_LOG_MAX_RETURN)) + with _WS_LOG_LOCK: + items = list(_WS_LOG_ENTRIES) + if len(items) <= cap: + return items + return items[-cap:] + + def parse_subscription_id_from_already_exists_error(message: str) -> str: """Extract subscription id from UME 400 'topic subscription already exist, id:...'.""" m = _ORPHAN_SUB_ID_RE.search(str(message or "")) @@ -107,7 +163,7 @@ def load_persisted_subscription() -> bool: return False sub_id, uri, topic = loaded _set_active_subscription(sub_id, uri, topic=topic) - _ws_log.info("loaded persisted alarm subscription id=%s", sub_id) + append_ws_log(f"loaded persisted subscription id={sub_id}", subscription_id=sub_id) return True finally: db.close() @@ -137,14 +193,14 @@ def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[ sub_id, uri, topic = existing _set_active_subscription(sub_id, uri, topic=topic) request_ws_reconnect() - _ws_log.info("establish skipped: subscription already exists id=%s", sub_id) + append_ws_log(f"establish skipped: already in DB id={sub_id}", subscription_id=sub_id) st = get_subscription_status() return {**st, "already_exists": True} mem_id, mem_uri = get_active_subscription() if mem_id and mem_uri: request_ws_reconnect() - _ws_log.info("establish skipped: in-memory subscription id=%s", mem_id) + append_ws_log(f"establish skipped: in-memory id={mem_id}", subscription_id=mem_id) st = get_subscription_status() return {**st, "already_exists": True} @@ -161,7 +217,11 @@ def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[ orphan_id = parse_subscription_id_from_already_exists_error(str(exc)) if not orphan_id: raise - _ws_log.warning("establish: UME reports existing subscription id=%s, deleting then retry", orphan_id) + append_ws_log( + f"establish: orphan on UME id={orphan_id}, deleting then retry", + level="warning", + subscription_id=orphan_id, + ) try: _delete_subscription_on_ume(client, orphan_id) except RuntimeError as del_exc: @@ -177,7 +237,7 @@ def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[ _active_client = client _set_active_subscription(sub_id, uri, topic=topic) request_ws_reconnect() - _ws_log.info("manual establish alarm subscription id=%s", sub_id) + append_ws_log(f"manual establish subscription id={sub_id}", subscription_id=sub_id) return {**get_subscription_status(), "already_exists": False} @@ -195,7 +255,7 @@ def cancel_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str db.commit() _clear_active_subscription() request_ws_reconnect() - _ws_log.info("manual cancel alarm subscription id=%s", sub_id) + append_ws_log(f"manual cancel subscription id={sub_id or '(none)'}", subscription_id=sub_id) return get_subscription_status() @@ -224,35 +284,48 @@ def _run_ws_session( def _on_message(_ws: Any, message: str) -> None: payload = _parse_ws_message(message) if payload is None: + append_ws_log("message ignored: invalid JSON", level="warning", subscription_id=subscription_id) return db = SessionLocal() try: - action, changed = process_alarm_notification(db, payload) + alarm = extract_alarm_from_notification(payload) + if alarm is None: + append_ws_log("收包忽略: 非告警通知", level="warning", subscription_id=subscription_id) + return + norm = normalize_yang_alarm(alarm) or {} + alarm_key = _alarm_key(norm) or "?" + action, changed = apply_alarm_to_current(db, alarm, touch_ts=_utc_now_naive()) if changed: db.commit() else: db.rollback() + 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() - _ws_log.exception("ws alarm apply failed: %s", exc) + 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: - _ws_log.warning("ws error subscription=%s: %s", subscription_id, error) + 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]}") def _on_close(_ws: Any, close_status_code: Any, close_msg: Any) -> None: - _ws_log.info("ws closed subscription=%s code=%s msg=%s", subscription_id, close_status_code, close_msg) + append_ws_log( + f"ws closed code={close_status_code} msg={str(close_msg or '')[:120]}", + subscription_id=subscription_id, + ) closed.set() def _on_open(_ws: Any) -> None: - _ws_log.info("ws connected subscription=%s", subscription_id) + append_ws_log(f"ws connected uri={wss_uri[:120]}", subscription_id=subscription_id) if on_status is not None: on_status("connected") @@ -317,28 +390,31 @@ 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) + while stop_event is None or not stop_event.is_set(): if is_paused is not None and is_paused(): - if on_status is not None: - on_status("paused") + _status("paused", level="warning") time.sleep(1.0) continue subscription_id, wss_uri = get_active_subscription() if not subscription_id or not wss_uri: - if on_status is not None: - on_status("no_subscription") + _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): - if on_status is not None: - on_status("waiting_token") + _status("waiting_token", level="warning", sub_id=subscription_id) time.sleep(2.0) continue + append_ws_log(f"connecting {wss_uri[:160]}", subscription_id=subscription_id) _run_ws_session( client, wss_uri=wss_uri, @@ -349,28 +425,21 @@ def run_alarm_ws_consumer_loop( backoff_s = 2.0 except RuntimeError as exc: if "ume_ws_no_valid_token" in str(exc): - if on_status is not None: - on_status("waiting_token") + _status("waiting_token", level="warning", sub_id=subscription_id) time.sleep(2.0) continue - _ws_log.exception("alarm ws session failed: %s", exc) - if on_status is not None: - on_status(f"error:{str(exc)[:120]}") + _status(f"session error: {str(exc)[:200]}", level="error", sub_id=subscription_id) except Exception as exc: - _ws_log.exception("alarm ws session failed: %s", exc) - if on_status is not None: - on_status(f"error:{str(exc)[:120]}") + _status(f"session error: {str(exc)[:200]}", level="error", sub_id=subscription_id) 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: - if on_status is not None: - on_status("no_subscription") + _status("no_subscription") continue - if on_status is not None: - on_status(f"reconnect_ws_in_{int(backoff_s)}s") + _status(f"reconnect in {int(backoff_s)}s", sub_id=sub_id_after) slept = 0.0 while slept < backoff_s: if stop_event is not None and stop_event.is_set(): @@ -402,9 +471,11 @@ def start_ume_alarm_ws_consumer( if not bool(getattr(settings, "ume_alarm_ws_enabled", True)): return if not str(client.base_url or "").strip(): + append_ws_log("consumer disabled: no UME base URL", level="warning") if on_status is not None: on_status("disabled:no_base_url") return + append_ws_log("ws consumer thread started") global _active_client with _shutdown_lock: _active_client = client diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index ef14fa8..1464f95 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -16,9 +16,11 @@ from netx_api import ume_alarm_ws from netx_api.models import UmeAlarmSubscription from netx_api.ume_alarm_subscription_store import clear_subscription, load_subscription, save_subscription from netx_api.ume_alarm_ws import ( + append_ws_log, cancel_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, + get_ws_logs, load_persisted_subscription, parse_subscription_id_from_already_exists_error, process_alarm_notification, @@ -736,6 +738,22 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertTrue(changed) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-1")) + def test_ws_log_ring_buffer(self): + from netx_api import ume_alarm_ws as ws_mod + + with ws_mod._WS_LOG_LOCK: + ws_mod._WS_LOG_ENTRIES.clear() + try: + for i in range(5): + append_ws_log(f"line-{i}") + logs = get_ws_logs(limit=3) + self.assertEqual(len(logs), 3) + self.assertEqual(logs[-1]["message"], "line-4") + self.assertIn("ts", logs[0]) + finally: + with ws_mod._WS_LOG_LOCK: + ws_mod._WS_LOG_ENTRIES.clear() + def test_parse_subscription_id_from_already_exists_error(self): msg = ( 'ume_request_failed:400:{"error":{"errorInfo":"topic subscription already exist, ' diff --git a/web/src/index.css b/web/src/index.css index 61dd515..fd99f46 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -189,6 +189,46 @@ pre { padding: 10px; } +.ws-log-panel { + margin-top: 12px; + max-height: 240px; + overflow-y: auto; + background: #091425; + border: 1px solid #223651; + border-radius: 8px; + padding: 8px 10px; + font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; + font-size: 12px; + line-height: 1.45; +} + +.ws-log-line { + padding: 3px 0; + border-bottom: 1px solid #152238; + word-break: break-word; +} + +.ws-log-line:last-child { + border-bottom: none; +} + +.ws-log-line--error { + color: #ff8a9a; +} + +.ws-log-line--warning { + color: #f0c060; +} + +.ws-log-line--info { + color: #b8c9e0; +} + +.ws-log-ts { + color: #6b8299; + margin-right: 8px; +} + .pill { margin-top: 8px; padding: 4px 8px; diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index b1f2b6d..da012c5 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -197,6 +197,7 @@ 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 wsLogs = [...(subscriptionStatusQuery.data?.ws_logs ?? alarmSub.ws_logs ?? [])].reverse(); const subscriptionActive = Boolean(alarmSub.active); const subPending = subscriptionEstablishMutation.isPending || subscriptionCancelMutation.isPending; @@ -364,6 +365,37 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { {subscriptionOpError ? (
订阅操作失败: {subscriptionOpError}
) : null} +
+
+ WSS 运行日志 + + {wsConsumer?.last_run_at + ? `最近活动 ${formatSystemTime(String(wsConsumer.last_run_at))}` + : "自动刷新"} + +
+
+ {wsLogs.length === 0 ? ( +
暂无日志(建立订阅并连接 WSS 后会出现连接、收包、重连等记录)
+ ) : ( + wsLogs.map((line, idx) => { + const lvl = String(line.level || "info").toLowerCase(); + const cls = + lvl === "error" ? "ws-log-line--error" : lvl === "warning" ? "ws-log-line--warning" : "ws-log-line--info"; + return ( +
+ {formatSystemTime(line.ts)} + [{lvl}] + {line.message} + {line.subscription_id ? ( + · {String(line.subscription_id).slice(0, 8)}… + ) : null} +
+ ); + }) + )} +
+

UME 同步

diff --git a/web/src/types.ts b/web/src/types.ts index dc8db9d..f25683c 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -89,6 +89,13 @@ export type UmeSyncJobItem = { ended_at?: string | null; }; +export type UmeWsLogEntry = { + ts: string; + level: string; + message: string; + subscription_id?: string; +}; + export type UmeAlarmSubscriptionStatus = { ok?: boolean; created?: boolean; @@ -100,6 +107,7 @@ export type UmeAlarmSubscriptionStatus = { ws_consumer_status?: string; ws_consumer_last_error?: string; ws_consumer_last_run_at?: string | null; + ws_logs?: UmeWsLogEntry[]; }; export type UmeSyncStatusResponse = {