diff --git a/netx_api/dsh_alarm_hub.py b/netx_api/dsh_alarm_hub.py index 233f9c9..b9dce6a 100644 --- a/netx_api/dsh_alarm_hub.py +++ b/netx_api/dsh_alarm_hub.py @@ -10,6 +10,8 @@ from __future__ import annotations import asyncio import logging import threading +import uuid +from dataclasses import dataclass, field from datetime import datetime, timezone from typing import Any @@ -24,7 +26,7 @@ from .db import SessionLocal _log = logging.getLogger("netx.dsh.alarm_hub") _LOCK = threading.Lock() -_CLIENTS: set[WebSocket] = set() +_CLIENTS: dict[WebSocket, "SubscriberInfo"] = {} _LOOP: asyncio.AbstractEventLoop | None = None _STATS = { "published": 0, @@ -34,10 +36,33 @@ _STATS = { } +@dataclass +class SubscriberInfo: + """One authenticated netxops (or other DSH) subscriber.""" + + id: str + user: str + remote: str = "" + client: str = "" + connected_at: str = field(default_factory=lambda: _utc_now_iso()) + last_seen_at: str = field(default_factory=lambda: _utc_now_iso()) + + def _utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat() +def _client_remote(websocket: WebSocket) -> str: + client = getattr(websocket, "client", None) + if client is None: + return "" + host = getattr(client, "host", None) or "" + port = getattr(client, "port", None) + if host and port is not None: + return f"{host}:{port}" + return str(host or "") + + def bind_event_loop(loop: asyncio.AbstractEventLoop | None) -> None: """Remember the API process loop so sync publishers can schedule sends.""" global _LOOP @@ -51,16 +76,28 @@ def subscriber_count() -> int: def hub_status() -> dict[str, Any]: with _LOCK: - clients = len(_CLIENTS) + clients = list(_CLIENTS.values()) stats = dict(_STATS) - stats["subscribers"] = clients + stats["subscribers"] = len(clients) + connections = [ + { + "id": info.id, + "user": info.user, + "remote": info.remote, + "client": info.client, + "connected_at": info.connected_at, + "last_seen_at": info.last_seen_at, + } + for info in sorted(clients, key=lambda x: x.connected_at) + ] return { "enabled": True, "path": "/v1/integrations/dsh-alarm/ws", - "subscribers": clients, + "subscribers": len(connections), "published": int(stats.get("published") or 0), "deliver_ok": int(stats.get("deliver_ok") or 0), "deliver_fail": int(stats.get("deliver_fail") or 0), + "connections": connections, } @@ -96,9 +133,19 @@ async def _send_json(ws: WebSocket, payload: dict[str, Any]) -> bool: return False +def _drop_clients(dead: list[WebSocket]) -> None: + if not dead: + return + with _LOCK: + for ws in dead: + _CLIENTS.pop(ws, None) + _STATS["subscribers"] = len(_CLIENTS) + _STATS["deliver_fail"] += len(dead) + + async def _broadcast(payload: dict[str, Any]) -> int: with _LOCK: - clients = list(_CLIENTS) + clients = list(_CLIENTS.keys()) if not clients: return 0 envelope = { @@ -115,11 +162,7 @@ async def _broadcast(payload: dict[str, Any]) -> int: else: dead.append(ws) if dead: - with _LOCK: - for ws in dead: - _CLIENTS.discard(ws) - _STATS["subscribers"] = len(_CLIENTS) - _STATS["deliver_fail"] += len(dead) + _drop_clients(dead) for ws in dead: try: await ws.close() @@ -161,6 +204,7 @@ async def dsh_alarm_ws_loop(websocket: WebSocket) -> None: await websocket.accept() bind_event_loop(asyncio.get_running_loop()) authed = False + info: SubscriberInfo | None = None try: while True: raw = await websocket.receive_text() @@ -185,20 +229,41 @@ async def dsh_alarm_ws_loop(websocket: WebSocket) -> None: await _send_json(websocket, {"type": "auth-fail", "error": detail}) await websocket.close(code=4401) return + client_label = str( + msg.get("client") or msg.get("host") or msg.get("client_id") or "" + ).strip()[:120] + now = _utc_now_iso() + info = SubscriberInfo( + id=uuid.uuid4().hex[:12], + user=detail, + remote=_client_remote(websocket), + client=client_label, + connected_at=now, + last_seen_at=now, + ) authed = True with _LOCK: - _CLIENTS.add(websocket) + _CLIENTS[websocket] = info _STATS["subscribers"] = len(_CLIENTS) await _send_json( websocket, { "type": "auth-ok", "user": detail, - "ts": _utc_now_iso(), + "connection_id": info.id, + "ts": now, }, ) - _log.info("dsh alarm hub subscriber connected (%s)", detail) + _log.info( + "dsh alarm hub subscriber connected id=%s user=%s remote=%s client=%s", + info.id, + detail, + info.remote, + info.client or "-", + ) continue + if info is not None: + info.last_seen_at = _utc_now_iso() if mtype == "ping": await _send_json(websocket, {"type": "pong", "ts": _utc_now_iso()}) continue @@ -207,6 +272,10 @@ async def dsh_alarm_ws_loop(websocket: WebSocket) -> None: return finally: with _LOCK: - _CLIENTS.discard(websocket) + _CLIENTS.pop(websocket, None) _STATS["subscribers"] = len(_CLIENTS) - _log.info("dsh alarm hub subscriber disconnected") + _log.info( + "dsh alarm hub subscriber disconnected id=%s user=%s", + info.id if info else "-", + info.user if info else "-", + ) diff --git a/netx_api/ume_key_alert_router.py b/netx_api/ume_key_alert_router.py index 6e2da2f..ffc333b 100644 --- a/netx_api/ume_key_alert_router.py +++ b/netx_api/ume_key_alert_router.py @@ -34,6 +34,7 @@ from .models import ( UmeKeyAlertRule, UmeSyncJob, ) +from .dsh_alarm_hub import hub_status as dsh_alarm_hub_status from .oclaw_alarm_forwarder import ( forwarder_status, request_forwarder_reconnect, @@ -146,7 +147,15 @@ def ume_list_key_alert_rules( for row in rows ] fwd = forwarder_status() - return {"items": items, "total": total, "page": page, "page_size": page_size, "forwarder": fwd} + hub = dsh_alarm_hub_status() + return { + "items": items, + "total": total, + "page": page, + "page_size": page_size, + "forwarder": fwd, + "dsh_alarm_hub": hub, + } @router.get("/v1/ume/key-alert-monitor") @@ -173,6 +182,8 @@ def ume_key_alert_monitor( "page": int(base.get("page") or page), "page_size": int(base.get("page_size") or page_size), "config": get_key_alert_monitor_config(db), + "dsh_alarm_hub": base.get("dsh_alarm_hub") or dsh_alarm_hub_status(), + # Legacy OClaw outbound bridge (optional; demoted in UI). "forwarder": base.get("forwarder") or forwarder_status(), } diff --git a/tests/test_dsh_alarm_hub.py b/tests/test_dsh_alarm_hub.py new file mode 100644 index 0000000..9265ae7 --- /dev/null +++ b/tests/test_dsh_alarm_hub.py @@ -0,0 +1,72 @@ +"""Unit tests for DSH alarm hub subscriber bookkeeping.""" + +from __future__ import annotations + +import unittest +from types import SimpleNamespace +from unittest.mock import MagicMock + +from netx_api import dsh_alarm_hub as hub + + +class DshAlarmHubStatusTests(unittest.TestCase): + def setUp(self) -> None: + with hub._LOCK: + hub._CLIENTS.clear() + hub._STATS.update( + { + "published": 0, + "deliver_ok": 0, + "deliver_fail": 0, + "subscribers": 0, + } + ) + + def tearDown(self) -> None: + with hub._LOCK: + hub._CLIENTS.clear() + + def test_hub_status_lists_connections(self) -> None: + ws_a = MagicMock(name="ws_a") + ws_b = MagicMock(name="ws_b") + info_a = hub.SubscriberInfo( + id="aaa111", + user="alice", + remote="10.0.0.1:40001", + client="netxops@host-a", + connected_at="2026-09-10T08:00:00+00:00", + last_seen_at="2026-09-10T08:01:00+00:00", + ) + info_b = hub.SubscriberInfo( + id="bbb222", + user="bob", + remote="10.0.0.2:40002", + client="netxops@host-b", + connected_at="2026-09-10T08:00:30+00:00", + last_seen_at="2026-09-10T08:01:30+00:00", + ) + with hub._LOCK: + hub._CLIENTS[ws_a] = info_a + hub._CLIENTS[ws_b] = info_b + hub._STATS["subscribers"] = 2 + hub._STATS["published"] = 5 + hub._STATS["deliver_ok"] = 8 + + status = hub.hub_status() + self.assertTrue(status["enabled"]) + self.assertEqual(status["subscribers"], 2) + self.assertEqual(status["published"], 5) + self.assertEqual(status["deliver_ok"], 8) + self.assertEqual(status["path"], "/v1/integrations/dsh-alarm/ws") + ids = [row["id"] for row in status["connections"]] + self.assertEqual(ids, ["aaa111", "bbb222"]) + self.assertEqual(status["connections"][0]["user"], "alice") + self.assertEqual(status["connections"][1]["client"], "netxops@host-b") + + def test_client_remote_formats_host_port(self) -> None: + ws = SimpleNamespace(client=SimpleNamespace(host="192.168.1.9", port=54321)) + self.assertEqual(hub._client_remote(ws), "192.168.1.9:54321") + + +if __name__ == "__main__": + unittest.main() diff --git a/web/WEB.md b/web/WEB.md index c00e803..c030633 100644 --- a/web/WEB.md +++ b/web/WEB.md @@ -81,6 +81,7 @@ src/ - Query key 集中在 `constants/queryKeys.ts` - 失效缓存用 prefix key(如 `queryKeys.umeSyncStatusAll`) - 顶栏连接状态:`App` 轮询 `GET /v1/integrations/status`(5s),展示 **netx api** 与 **oclaw bridge**(含延迟 / 错误类型) +- 关键告警推送主路径:NetX **DSH alarm hub**(`/v1/integrations/dsh-alarm/ws`)为服务器,Netx Ops 外拨订阅;UME 页展示多订阅连接列表。OClaw 出站 WSS 为遗留旁路(`NETX_OCLAW_ALARM_WS_ENABLED`) ## 网元管理(独立于 UME) diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 9f4b65a..f6403ab 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -1025,7 +1025,24 @@ const en = { title: "AI alarm monitor", showPanel: "Details", hidePanel: "Close", - help: "Exact match by notificationId, or substring match on alarm description (case-insensitive; nativeProbableCause, etc.). Optionally restrict to UME inventory ne_type values (multi-select; empty = all devices). Pushed via OClaw WebSocket. Shares NETX_OCLAW_ANALYZE_TOKEN.", + help: "Exact match by notificationId, or substring match on alarm description (case-insensitive; nativeProbableCause, etc.). Optionally restrict to UME inventory ne_type values (multi-select; empty = all devices). Matched alerts are broadcast by NetX (server) to subscribed Netx Ops (DSH) clients.", + hub: "DSH subscribers", + hubConnected: "subscribed", + hubEmpty: "no subscribers", + hubSubscribers: "subscribers", + hubPublished: "published", + hubDeliverOk: "deliver ok", + hubDeliverFail: "deliver fail", + hubConnections: "Active connections", + hubColId: "id", + hubColUser: "user", + hubColRemote: "remote", + hubColClient: "client", + hubColConnected: "connected", + hubColLastSeen: "last seen", + hubEmptyConnections: "No Netx Ops clients are connected to the alarm hub", + hubPath: "path", + legacyOclaw: "Legacy OClaw outbound", ws: "OClaw WSS", wsConnected: "connected", wsDisconnected: "disconnected", @@ -1046,7 +1063,7 @@ const en = { neTypesEmpty: "No device types yet — sync Inventory first", neTypesAll: "All", forwardOnClear: "Push on alarm clear", - forwardOnClearHelp: "Global switch for all monitor rules. By default only new/updated alarms are pushed; when enabled, clear notifications are also sent when UME reports alarms cleared.", + forwardOnClearHelp: "Global switch for all monitor rules. By default only new/updated alarms are pushed; when enabled, clear notifications are also sent to subscribed Netx Ops clients.", add: "Add rules", adding: "Adding…", delete: "Delete", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index ddeb95c..795c03d 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -1017,7 +1017,24 @@ const zh = { title: "AI 告警监控", showPanel: "详情", hidePanel: "关闭", - help: "按 notificationId 精确匹配,或按告警描述关键字模糊匹配(不区分大小写,匹配 nativeProbableCause 等字段);可选限定 UME Inventory 中的网元类型 (ne_type),不选则匹配全部设备。经 OClaw WebSocket 推送到 WhatsApp 群。鉴权与 NETX_OCLAW_ANALYZE_TOKEN 共用。", + help: "按 notificationId 精确匹配,或按告警描述关键字模糊匹配(不区分大小写,匹配 nativeProbableCause 等字段);可选限定 UME Inventory 中的网元类型 (ne_type),不选则匹配全部设备。命中后由 NetX 作为服务器广播给已订阅的 Netx Ops(DSH)客户端。", + hub: "DSH 订阅", + hubConnected: "有订阅", + hubEmpty: "无订阅", + hubSubscribers: "订阅数", + hubPublished: "已广播", + hubDeliverOk: "投递成功", + hubDeliverFail: "投递失败", + hubConnections: "当前连接", + hubColId: "连接 ID", + hubColUser: "用户", + hubColRemote: "远端", + hubColClient: "客户端", + hubColConnected: "连入时间", + hubColLastSeen: "最近心跳", + hubEmptyConnections: "当前没有 Netx Ops 客户端连入告警 hub", + hubPath: "路径", + legacyOclaw: "遗留 OClaw 出站", ws: "OClaw WSS", wsConnected: "已连接", wsDisconnected: "未连接", @@ -1038,7 +1055,7 @@ const zh = { neTypesEmpty: "暂无设备类型,请先同步 Inventory", neTypesAll: "全部", forwardOnClear: "告警清除时也推送", - forwardOnClearHelp: "全局开关:对所有监控规则生效。默认只在告警新增/更新时推送;勾选后,UME 上报告警已清除时也会向 WhatsApp 发送清除通知。", + forwardOnClearHelp: "全局开关:对所有监控规则生效。默认只在告警新增/更新时推送;勾选后,UME 上报告警已清除时也会向已订阅的 Netx Ops 推送清除通知。", add: "添加规则", adding: "添加中…", delete: "删除", diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index 5d5994b..267f4d4 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -432,11 +432,15 @@ export function UmePage() { [neQuery.data, expandedNeId], ); + const keyAlertHub = keyAlertMonitorQuery.data?.dsh_alarm_hub; + const keyAlertHubConnections = keyAlertHub?.connections || []; + const keyAlertHubSubscribers = Number(keyAlertHub?.subscribers || keyAlertHubConnections.length || 0); const keyAlertForwarder = keyAlertMonitorQuery.data?.forwarder; const keyAlertRules = keyAlertMonitorQuery.data?.rules || []; const keyAlertTotal = Number(keyAlertMonitorQuery.data?.total || keyAlertRules.length); const keyAlertPages = pageCount(keyAlertTotal, keyAlertPageSize); const keyAlertForwardOnClear = Boolean(keyAlertMonitorQuery.data?.config?.forward_on_clear); + const hubWsPill = keyAlertHubSubscribers > 0 ? "up" : "down"; const oclawWsPill = !keyAlertForwarder?.enabled ? "unknown" @@ -445,6 +449,7 @@ export function UmePage() { : keyAlertForwarder.connected ? "up" : "down"; + const showLegacyOclaw = Boolean(keyAlertForwarder?.enabled); const runtimeTaskLabel = (task: string) => { const key = `ume.tasks.runtimeTask.${task}`; @@ -852,25 +857,30 @@ export function UmePage() {
- - {t("ume.keyAlert.ws")}:{" "} - {!keyAlertForwarder?.enabled - ? t("ume.keyAlert.wsDisabled") - : keyAlertForwarder.paused + + {t("ume.keyAlert.hub")}:{" "} + {keyAlertHubSubscribers > 0 + ? t("ume.keyAlert.hubConnected") + : t("ume.keyAlert.hubEmpty")} + + + {t("ume.keyAlert.hubSubscribers")}: {keyAlertHubSubscribers} + + + {t("ume.keyAlert.hubPublished")}: {Number(keyAlertHub?.published || 0)} + + + {t("ume.keyAlert.hubDeliverOk")}: {Number(keyAlertHub?.deliver_ok || 0)} + + {showLegacyOclaw ? ( + + {t("ume.keyAlert.ws")}:{" "} + {keyAlertForwarder?.paused ? t("ume.keyAlert.wsPaused") - : keyAlertForwarder.connected + : keyAlertForwarder?.connected ? t("ume.keyAlert.wsConnected") : t("ume.keyAlert.wsDisconnected")} - - {keyAlertForwarder?.enabled ? ( - <> - - {t("ume.keyAlert.publishedOk")}: {Number(keyAlertForwarder.published_ok || 0)} - - - {t("ume.keyAlert.queue")}: {Number(keyAlertForwarder.queue_size || 0)} - - + ) : null}
@@ -964,16 +974,22 @@ export function UmePage() {

{t("ume.keyAlert.title")}

- {t("ume.keyAlert.ws")}:{" "} - {!keyAlertForwarder?.enabled - ? t("ume.keyAlert.wsDisabled") - : keyAlertForwarder.paused - ? t("ume.keyAlert.wsPaused") - : keyAlertForwarder.connected - ? t("ume.keyAlert.wsConnected") - : t("ume.keyAlert.wsDisconnected")} - {keyAlertForwarder?.enabled - ? ` · ${t("ume.keyAlert.publishedOk")} ${Number(keyAlertForwarder.published_ok || 0)} · ${t("ume.keyAlert.queue")} ${Number(keyAlertForwarder.queue_size || 0)}` + {t("ume.keyAlert.hub")}:{" "} + {keyAlertHubSubscribers > 0 + ? t("ume.keyAlert.hubConnected") + : t("ume.keyAlert.hubEmpty")} + {` · ${t("ume.keyAlert.hubSubscribers")} ${keyAlertHubSubscribers}`} + {` · ${t("ume.keyAlert.hubPublished")} ${Number(keyAlertHub?.published || 0)}`} + {` · ${t("ume.keyAlert.hubDeliverOk")} ${Number(keyAlertHub?.deliver_ok || 0)}`} + {keyAlertHub?.path ? ` · ${t("ume.keyAlert.hubPath")} ${keyAlertHub.path}` : ""} + {showLegacyOclaw + ? ` · ${t("ume.keyAlert.legacyOclaw")}: ${ + keyAlertForwarder?.paused + ? t("ume.keyAlert.wsPaused") + : keyAlertForwarder?.connected + ? t("ume.keyAlert.wsConnected") + : t("ume.keyAlert.wsDisconnected") + }` : ""}

@@ -992,6 +1008,62 @@ export function UmePage() {
+
+
{t("ume.keyAlert.hubConnections")}
+ {keyAlertHubConnections.length === 0 ? ( +

{t("ume.keyAlert.hubEmptyConnections")}

+ ) : ( +
+ + + + + + + + + + + + + {keyAlertHubConnections.map((conn) => ( + + + + + + + + + ))} + +
{t("ume.keyAlert.hubColId")}{t("ume.keyAlert.hubColUser")}{t("ume.keyAlert.hubColRemote")}{t("ume.keyAlert.hubColClient")}{t("ume.keyAlert.hubColConnected")}{t("ume.keyAlert.hubColLastSeen")}
+ {conn.id} + {conn.user || "-"}{conn.remote || "-"}{conn.client || "-"} + {conn.connected_at + ? formatSystemTime(conn.connected_at) + : "-"} + + {conn.last_seen_at + ? formatSystemTime(conn.last_seen_at) + : "-"} +
+
+ )} + {showLegacyOclaw ? ( +

+ {t("ume.keyAlert.legacyOclaw")}:{" "} + {keyAlertForwarder?.paused + ? t("ume.keyAlert.wsPaused") + : keyAlertForwarder?.connected + ? t("ume.keyAlert.wsConnected") + : t("ume.keyAlert.wsDisconnected")} + {` · ${t("ume.keyAlert.publishedOk")} ${Number(keyAlertForwarder?.published_ok || 0)}`} + {` · ${t("ume.keyAlert.queue")} ${Number(keyAlertForwarder?.queue_size || 0)}`} +

+ ) : null} +
+
diff --git a/web/src/types.ts b/web/src/types.ts index 32deb2a..2937f8e 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -1,6 +1,26 @@ +export type DshAlarmHubConnection = { + id: string; + user: string; + remote?: string; + client?: string; + connected_at?: string; + last_seen_at?: string; +}; + +export type DshAlarmHubStatus = { + enabled: boolean; + path?: string; + subscribers: number; + published?: number; + deliver_ok?: number; + deliver_fail?: number; + connections?: DshAlarmHubConnection[]; +}; + export type IntegrationStatus = { netx_api: { status: "up" | "down" | "unknown"; [k: string]: unknown }; db: { status: "up" | "down" | "unknown"; latency_ms?: number; error?: string; [k: string]: unknown }; + dsh_alarm_hub?: DshAlarmHubStatus; oclaw_bridge?: { status: "up" | "down" | "unknown"; mode?: string; @@ -58,6 +78,9 @@ export type UmeKeyAlertMonitorResponse = { config?: { forward_on_clear: boolean; }; + /** Primary: NetX hub ← netxops clients (multi-subscriber). */ + dsh_alarm_hub?: DshAlarmHubStatus; + /** Legacy: NetX → OClaw outbound bridge (single link). */ forwarder: UmeKeyAlertForwarderStatus; };