Show DSH alarm hub multi-subscriber status instead of OClaw single link.

Track per-connection metadata on the hub and surface it in the UME key-alert UI so operators can see which Netx Ops clients are subscribed.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-10 16:57:42 +08:00
parent 0f8a06bfba
commit 5096117ad4
8 changed files with 328 additions and 46 deletions

View file

@ -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 "-",
)

View file

@ -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(),
}

View file

@ -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()

View file

@ -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)

View file

@ -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",

View file

@ -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: "删除",

View file

@ -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() {
</button>
</div>
<div className="ume-entry__pills actions-row actions-row--inline">
<span className={`conn-pill conn-pill--${oclawWsPill}`}>
<span className={`conn-pill conn-pill--${hubWsPill}`}>
{t("ume.keyAlert.hub")}:{" "}
{keyAlertHubSubscribers > 0
? t("ume.keyAlert.hubConnected")
: t("ume.keyAlert.hubEmpty")}
</span>
<span className="conn-pill">
{t("ume.keyAlert.hubSubscribers")}: {keyAlertHubSubscribers}
</span>
<span className="conn-pill">
{t("ume.keyAlert.hubPublished")}: {Number(keyAlertHub?.published || 0)}
</span>
<span className="conn-pill">
{t("ume.keyAlert.hubDeliverOk")}: {Number(keyAlertHub?.deliver_ok || 0)}
</span>
{showLegacyOclaw ? (
<span className={`conn-pill conn-pill--${oclawWsPill}`} title={t("ume.keyAlert.legacyOclaw")}>
{t("ume.keyAlert.ws")}:{" "}
{!keyAlertForwarder?.enabled
? t("ume.keyAlert.wsDisabled")
: keyAlertForwarder.paused
{keyAlertForwarder?.paused
? t("ume.keyAlert.wsPaused")
: keyAlertForwarder.connected
: keyAlertForwarder?.connected
? t("ume.keyAlert.wsConnected")
: t("ume.keyAlert.wsDisconnected")}
</span>
{keyAlertForwarder?.enabled ? (
<>
<span className="conn-pill">
{t("ume.keyAlert.publishedOk")}: {Number(keyAlertForwarder.published_ok || 0)}
</span>
<span className="conn-pill">
{t("ume.keyAlert.queue")}: {Number(keyAlertForwarder.queue_size || 0)}
</span>
</>
) : null}
</div>
</article>
@ -964,16 +974,22 @@ export function UmePage() {
<div className="ops-detail-modal__title">
<h3>{t("ume.keyAlert.title")}</h3>
<p className="muted">
{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)}`}
{keyAlertHub?.path ? ` · ${t("ume.keyAlert.hubPath")} ${keyAlertHub.path}` : ""}
{showLegacyOclaw
? ` · ${t("ume.keyAlert.legacyOclaw")}: ${
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)}`
: t("ume.keyAlert.wsDisconnected")
}`
: ""}
</p>
</div>
@ -992,6 +1008,62 @@ export function UmePage() {
</div>
<div className="ops-detail-modal__scroll ops-detail-modal__scroll--flow">
<div className="ume-modal-form">
<div className="muted ume-modal-form__label">{t("ume.keyAlert.hubConnections")}</div>
{keyAlertHubConnections.length === 0 ? (
<p className="muted">{t("ume.keyAlert.hubEmptyConnections")}</p>
) : (
<div className="table-wrap">
<table className="data-table">
<thead>
<tr>
<th>{t("ume.keyAlert.hubColId")}</th>
<th>{t("ume.keyAlert.hubColUser")}</th>
<th>{t("ume.keyAlert.hubColRemote")}</th>
<th>{t("ume.keyAlert.hubColClient")}</th>
<th>{t("ume.keyAlert.hubColConnected")}</th>
<th>{t("ume.keyAlert.hubColLastSeen")}</th>
</tr>
</thead>
<tbody>
{keyAlertHubConnections.map((conn) => (
<tr key={conn.id}>
<td>
<code>{conn.id}</code>
</td>
<td>{conn.user || "-"}</td>
<td>{conn.remote || "-"}</td>
<td>{conn.client || "-"}</td>
<td>
{conn.connected_at
? formatSystemTime(conn.connected_at)
: "-"}
</td>
<td>
{conn.last_seen_at
? formatSystemTime(conn.last_seen_at)
: "-"}
</td>
</tr>
))}
</tbody>
</table>
</div>
)}
{showLegacyOclaw ? (
<p className="muted ume-modal-form__hint">
{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)}`}
</p>
) : null}
</div>
<div className="ume-modal-form">
<div className="ops-detail-modal__toolbar filter-inline">
<label className="muted">{t("ume.keyAlert.matchType")}</label>

View file

@ -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;
};