feat(ume): real-time current alarms via WebSocket subscription

Add WS consumer, persisted subscription store, manual subscribe/cancel APIs, and UI controls; extend sync and tests.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-22 20:27:43 +08:00
parent 21f7cfba41
commit 4f8735a83f
13 changed files with 1171 additions and 22 deletions

View file

@ -20,4 +20,8 @@ NETX_UME_TOKEN_LOGOUT_PATH=/restconf/operations/zte-security:oauth_token
NETX_UME_NE_PATH=/restconf/data/zte-resources-module:network-elements NETX_UME_NE_PATH=/restconf/data/zte-resources-module:network-elements
NETX_UME_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list NETX_UME_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list
NETX_UME_SYNC_ALARMS_CURRENT_ENABLED=true NETX_UME_SYNC_ALARMS_CURRENT_ENABLED=true
NETX_UME_SYNC_ALARMS_CURRENT_INTERVAL_S=300 NETX_UME_SYNC_ALARMS_CURRENT_INTERVAL_S=18000
NETX_UME_ALARM_WS_ENABLED=true
NETX_UME_NOTIFICATION_ESTABLISH_PATH=/restconf/operations/zte-notifications:establish-subscription
NETX_UME_NOTIFICATION_DELETE_PATH=/restconf/operations/zte-notifications:delete-subscription
NETX_UME_NOTIFICATION_TOPIC=ALARM

View file

@ -41,7 +41,11 @@ class Settings(BaseSettings):
ume_keepalive_interval_s: int = 600 ume_keepalive_interval_s: int = 600
ume_keepalive_renew_before_s: int = 900 ume_keepalive_renew_before_s: int = 900
ume_sync_alarms_current_enabled: bool = True ume_sync_alarms_current_enabled: bool = True
ume_sync_alarms_current_interval_s: int = 300 ume_sync_alarms_current_interval_s: int = 18000
ume_alarm_ws_enabled: bool = True
ume_notification_establish_path: str = "/restconf/operations/zte-notifications:establish-subscription"
ume_notification_delete_path: str = "/restconf/operations/zte-notifications:delete-subscription"
ume_notification_topic: str = "ALARM"
ume_sync_inventory_auto_enabled: bool = True ume_sync_inventory_auto_enabled: bool = True
ume_sync_inventory_every_hours: int = 48 ume_sync_inventory_every_hours: int = 48
ume_token_path: str = "/restconf/operations/zte-security:oauth_token" ume_token_path: str = "/restconf/operations/zte-security:oauth_token"

View file

@ -34,6 +34,14 @@ from .models import (
from .models import ImportJob from .models import ImportJob
from .parser_config import load_parser_config from .parser_config import load_parser_config
from .ume_client import UMEClient from .ume_client import UMEClient
from .ume_alarm_ws import (
cancel_alarm_subscription_manual,
establish_alarm_subscription_manual,
get_subscription_status,
load_persisted_subscription,
shutdown_ws_consumer,
start_ume_alarm_ws_consumer,
)
from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full
from .ume_token_store import ( from .ume_token_store import (
clear_shared_token, clear_shared_token,
@ -73,8 +81,10 @@ _SQL_FORBIDDEN_RE = re.compile(
_UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = { _UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = {
"token_keepalive": {"task": "token_keepalive", "status": "init", "last_run_at": None, "last_error": ""}, "token_keepalive": {"task": "token_keepalive", "status": "init", "last_run_at": None, "last_error": ""},
"alarms_current_auto_sync": {"task": "alarms_current_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, "alarms_current_auto_sync": {"task": "alarms_current_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
"alarms_current_ws_consumer": {"task": "alarms_current_ws_consumer", "status": "init", "last_run_at": None, "last_error": ""},
"inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, "inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
} }
_UME_WS_STOP_EVENT: threading.Event | None = None
_UME_RUNTIME_PAUSED: dict[str, bool] = {} _UME_RUNTIME_PAUSED: dict[str, bool] = {}
UME_KNOWN_RUNTIME_TASKS: tuple[str, ...] = tuple(_UME_RUNTIME_TASKS.keys()) UME_KNOWN_RUNTIME_TASKS: tuple[str, ...] = tuple(_UME_RUNTIME_TASKS.keys())
_UME_RUNTIME_LOCK = threading.Lock() _UME_RUNTIME_LOCK = threading.Lock()
@ -174,9 +184,13 @@ def _runtime_task_interval_fields(task_id: str) -> tuple[int | None, str]:
if task_id == "alarms_current_auto_sync": if task_id == "alarms_current_auto_sync":
if not bool(getattr(settings, "ume_sync_alarms_current_enabled", True)): if not bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
return None, "未启用" return None, "未启用"
interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 300) or 300) interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000)
eff = max(30, min(interval_s, 86400)) eff = max(30, min(interval_s, 86400))
return eff, _format_runtime_interval_label(eff) return eff, _format_runtime_interval_label(eff)
if task_id == "alarms_current_ws_consumer":
if not bool(getattr(settings, "ume_alarm_ws_enabled", True)):
return None, "未启用"
return None, "实时"
if task_id == "inventory_auto_sync": if task_id == "inventory_auto_sync":
if not bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)): if not bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)):
return None, "未启用" return None, "未启用"
@ -701,7 +715,7 @@ def on_startup() -> None:
) )
try: try:
if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)): if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 300) or 300) alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000)
alarms_interval_s = max(30, min(alarms_interval_s, 86400)) alarms_interval_s = max(30, min(alarms_interval_s, 86400))
def _alarms_current_sync_loop() -> None: def _alarms_current_sync_loop() -> None:
@ -832,6 +846,46 @@ def on_startup() -> None:
last_run_at=datetime.now(timezone.utc), last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}", last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
) )
global _UME_WS_STOP_EVENT
try:
if bool(getattr(settings, "ume_alarm_ws_enabled", True)) and str(getattr(settings, "ume_base_url", "") or "").strip():
if load_persisted_subscription():
_schedule_log.info("startup: loaded persisted UME alarm subscription")
_UME_WS_STOP_EVENT = threading.Event()
def _ws_on_status(msg: str) -> None:
_set_runtime_task(
"alarms_current_ws_consumer",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=str(msg or "")[:240],
)
t_ws = start_ume_alarm_ws_consumer(
_ume_client(),
on_status=_ws_on_status,
stop_event=_UME_WS_STOP_EVENT,
is_paused=lambda: _runtime_is_paused("alarms_current_ws_consumer"),
)
_schedule_log.info("started thread %s alive=%s", t_ws.name, t_ws.is_alive())
else:
_set_runtime_task("alarms_current_ws_consumer", status="paused", last_error="未启用或未配置 UME_BASE_URL")
except Exception as exc:
_schedule_log.exception("startup: alarms_current_ws_consumer thread init failed: %s", exc)
_set_runtime_task(
"alarms_current_ws_consumer",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
@app.on_event("shutdown")
def on_shutdown() -> None:
global _UME_WS_STOP_EVENT
if _UME_WS_STOP_EVENT is not None:
_UME_WS_STOP_EVENT.set()
shutdown_ws_consumer()
@app.get("/health") @app.get("/health")
@ -872,6 +926,41 @@ def ume_token_disconnect() -> dict[str, Any]:
return {"ok": ok, **st} return {"ok": ok, **st}
@app.get("/v1/ume/alarm-subscription/status")
def ume_alarm_subscription_status() -> dict[str, Any]:
st = get_subscription_status()
ws_task = _UME_RUNTIME_TASKS.get("alarms_current_ws_consumer") or {}
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"),
}
@app.post("/v1/ume/alarm-subscription/establish")
def ume_alarm_subscription_establish(db: Session = Depends(get_db)) -> dict[str, Any]:
client = _ume_client()
try:
st = establish_alarm_subscription_manual(client, db)
return {"ok": True, "created": not bool(st.get("already_exists")), **st}
except Exception as exc:
msg = str(exc)[:240]
raise HTTPException(status_code=502, detail=msg) from exc
@app.post("/v1/ume/alarm-subscription/cancel")
def ume_alarm_subscription_cancel(db: Session = Depends(get_db)) -> dict[str, Any]:
client = _ume_client()
try:
st = cancel_alarm_subscription_manual(client, db)
return {"ok": True, **st}
except Exception as exc:
msg = str(exc)[:240]
raise HTTPException(status_code=502, detail=msg) from exc
@app.post("/v1/ume/sync") @app.post("/v1/ume/sync")
def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict[str, Any]: def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict[str, Any]:
body = payload or {} body = payload or {}
@ -993,6 +1082,7 @@ def ume_sync_status(
"items": items, "items": items,
"latest_by_domain": latest_by_domain, "latest_by_domain": latest_by_domain,
"runtime_tasks": _list_runtime_tasks(), "runtime_tasks": _list_runtime_tasks(),
"alarm_subscription": get_subscription_status(),
} }

View file

@ -199,6 +199,19 @@ class UmeAlarmHistory(Base):
raw_json: Mapped[str] = mapped_column(Text, default="{}") raw_json: Mapped[str] = mapped_column(Text, default="{}")
class UmeAlarmSubscription(Base):
"""Persisted UME ALARM notification subscription (manual establish/cancel)."""
__tablename__ = "ume_alarm_subscription"
cache_key: Mapped[str] = mapped_column(String(64), primary_key=True, default="default")
subscription_id: Mapped[str] = mapped_column(String(128), default="", index=True)
wss_uri: Mapped[str] = mapped_column(Text, default="")
topic: Mapped[str] = mapped_column(String(64), default="ALARM")
established_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
class UmeTokenCache(Base): class UmeTokenCache(Base):
__tablename__ = "ume_token_cache" __tablename__ = "ume_token_cache"

View file

@ -0,0 +1,62 @@
from __future__ import annotations
from datetime import datetime, timezone
from sqlalchemy.orm import Session
from .models import UmeAlarmSubscription
DEFAULT_SUBSCRIPTION_KEY = "default"
def _utc_now_naive() -> datetime:
return datetime.now(timezone.utc).replace(tzinfo=None)
def load_subscription(db: Session, *, cache_key: str = DEFAULT_SUBSCRIPTION_KEY) -> tuple[str, str, str] | None:
row = db.get(UmeAlarmSubscription, str(cache_key or DEFAULT_SUBSCRIPTION_KEY))
if row is None:
return None
sub_id = str(row.subscription_id or "").strip()
uri = str(row.wss_uri or "").strip()
if not sub_id or not uri:
return None
return sub_id, uri, str(row.topic or "ALARM")
def save_subscription(
db: Session,
*,
subscription_id: str,
wss_uri: str,
topic: str = "ALARM",
cache_key: str = DEFAULT_SUBSCRIPTION_KEY,
) -> UmeAlarmSubscription:
key = str(cache_key or DEFAULT_SUBSCRIPTION_KEY)
now = _utc_now_naive()
row = db.get(UmeAlarmSubscription, key)
if row is None:
row = UmeAlarmSubscription(
cache_key=key,
subscription_id=str(subscription_id or "").strip(),
wss_uri=str(wss_uri or "").strip(),
topic=str(topic or "ALARM").strip() or "ALARM",
established_at=now,
updated_at=now,
)
db.add(row)
else:
row.subscription_id = str(subscription_id or "").strip()
row.wss_uri = str(wss_uri or "").strip()
row.topic = str(topic or "ALARM").strip() or "ALARM"
row.updated_at = now
db.flush()
return row
def clear_subscription(db: Session, *, cache_key: str = DEFAULT_SUBSCRIPTION_KEY) -> None:
key = str(cache_key or DEFAULT_SUBSCRIPTION_KEY)
row = db.get(UmeAlarmSubscription, key)
if row is not None:
db.delete(row)
db.flush()

402
netx_api/ume_alarm_ws.py Normal file
View file

@ -0,0 +1,402 @@
from __future__ import annotations
import json
import logging
import ssl
import threading
import time
from typing import Any, Callable
from sqlalchemy.orm import Session
from .config import settings
from .db import SessionLocal
from .ume_alarm_subscription_store import (
DEFAULT_SUBSCRIPTION_KEY,
clear_subscription,
load_subscription,
save_subscription,
)
from .ume_client import UMEClient
from .ume_sync_service import apply_alarm_to_current, extract_alarm_from_notification, _utc_now_naive
_ws_log = logging.getLogger("netx.ume.alarm_ws")
_shutdown_lock = threading.Lock()
_subscription_lock = threading.Lock()
_ws_wake_event = threading.Event()
_subscription_id: str = ""
_subscription_uri: str = ""
_subscription_topic: str = "ALARM"
_active_client: UMEClient | None = None
def _parse_ws_message(raw: str) -> dict[str, Any] | None:
text = str(raw or "").strip()
if not text:
return None
try:
data = json.loads(text)
except json.JSONDecodeError:
return None
return data if isinstance(data, dict) else None
def process_alarm_notification(db: Session, payload: dict[str, Any]) -> tuple[str, bool]:
alarm = extract_alarm_from_notification(payload)
if alarm is None:
return "skipped", False
return apply_alarm_to_current(db, alarm, touch_ts=_utc_now_naive())
def _delete_subscription_on_ume(client: UMEClient, subscription_id: str) -> None:
sub_id = str(subscription_id or "").strip()
if not sub_id:
return
client.delete_alarm_subscription(sub_id)
def _wait_for_shared_token(client: UMEClient, *, timeout_s: float = 300.0) -> bool:
deadline = time.time() + max(5.0, float(timeout_s))
while time.time() < deadline:
client._sync_token_from_store()
if client.has_valid_token():
return True
time.sleep(2.0)
return False
def _set_active_subscription(subscription_id: str, wss_uri: str, *, topic: str = "ALARM") -> None:
global _subscription_id, _subscription_uri, _subscription_topic
with _subscription_lock:
_subscription_id = str(subscription_id or "").strip()
_subscription_uri = str(wss_uri or "").strip()
_subscription_topic = str(topic or "ALARM").strip() or "ALARM"
def _clear_active_subscription() -> None:
global _subscription_id, _subscription_uri, _subscription_topic
with _subscription_lock:
_subscription_id = ""
_subscription_uri = ""
_subscription_topic = "ALARM"
def get_active_subscription() -> tuple[str, str]:
with _subscription_lock:
return str(_subscription_id or ""), str(_subscription_uri or "")
def load_persisted_subscription() -> bool:
"""Load subscription from DB into memory (call on process startup)."""
db = SessionLocal()
try:
loaded = load_subscription(db)
if loaded is None:
_clear_active_subscription()
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)
return True
finally:
db.close()
def get_subscription_status() -> dict[str, Any]:
with _subscription_lock:
sub_id = str(_subscription_id or "")
uri = str(_subscription_uri or "")
topic = str(_subscription_topic or "ALARM")
return {
"active": bool(sub_id and uri),
"subscription_id": sub_id,
"wss_uri": uri,
"topic": topic,
}
def request_ws_reconnect() -> None:
_ws_wake_event.set()
def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str, Any]:
"""Manually establish ALARM subscription (UME API + persist). Does not open WSS."""
existing = load_subscription(db)
if existing is not None:
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)
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)
st = get_subscription_status()
return {**st, "already_exists": True}
if not _wait_for_shared_token(client, timeout_s=120.0):
raise RuntimeError("ume_ws_no_valid_token:wait_timeout")
client._sync_token_from_store()
if not client.has_valid_token():
raise RuntimeError("ume_ws_no_valid_token")
topic = str(getattr(settings, "ume_notification_topic", "ALARM") or "ALARM").strip() or "ALARM"
sub_id, uri = client.establish_alarm_subscription(topic=topic)
save_subscription(db, subscription_id=sub_id, wss_uri=uri, topic=topic)
db.commit()
global _active_client
with _shutdown_lock:
_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)
return {**get_subscription_status(), "already_exists": False}
def cancel_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str, Any]:
"""Manually cancel subscription (UME delete + clear DB + drop WSS)."""
loaded = load_subscription(db)
sub_id, _uri = get_active_subscription()
if not sub_id and loaded is not None:
sub_id = str(loaded[0] or "")
if sub_id:
try:
client._sync_token_from_store()
if client.has_valid_token():
_delete_subscription_on_ume(client, sub_id)
except Exception as exc:
_ws_log.warning("cancel subscription on UME failed: %s", exc)
clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY)
db.commit()
_clear_active_subscription()
request_ws_reconnect()
_ws_log.info("manual cancel alarm subscription id=%s", sub_id)
return get_subscription_status()
def _run_ws_session(
client: UMEClient,
*,
wss_uri: str,
subscription_id: str,
on_status: Callable[[str], None] | None = None,
stop_event: threading.Event | None = None,
) -> None:
import websocket
headers = client.ws_auth_headers()
header_list = [f"{k}: {v}" for k, v in headers.items() if str(v or "").strip()]
sslopt: dict[str, Any] | None = None
if wss_uri.lower().startswith("wss://"):
if client.verify_tls:
sslopt = {"cert_reqs": ssl.CERT_REQUIRED}
else:
sslopt = {"cert_reqs": ssl.CERT_NONE}
closed = threading.Event()
def _on_message(_ws: Any, message: str) -> None:
payload = _parse_ws_message(message)
if payload is None:
return
db = SessionLocal()
try:
action, changed = process_alarm_notification(db, payload)
if changed:
db.commit()
else:
db.rollback()
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)
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)
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)
closed.set()
def _on_open(_ws: Any) -> None:
_ws_log.info("ws connected subscription=%s", subscription_id)
if on_status is not None:
on_status("connected")
ws_app = websocket.WebSocketApp(
wss_uri,
header=header_list,
on_open=_on_open,
on_message=_on_message,
on_error=_on_error,
on_close=_on_close,
)
kwargs: dict[str, Any] = {"ping_interval": 30, "ping_timeout": 20}
if sslopt is not None:
kwargs["sslopt"] = sslopt
thread = threading.Thread(
target=lambda: ws_app.run_forever(**kwargs),
name=f"ume-alarm-ws-{subscription_id[:8]}",
daemon=True,
)
thread.start()
while thread.is_alive() and not closed.is_set():
if stop_event is not None and stop_event.is_set():
try:
ws_app.close()
except Exception:
pass
break
sub_id_now, _ = get_active_subscription()
if not sub_id_now or sub_id_now != subscription_id:
try:
ws_app.close()
except Exception:
pass
break
if _ws_wake_event.is_set():
try:
ws_app.close()
except Exception:
pass
break
time.sleep(0.5)
if thread.is_alive():
thread.join(timeout=5.0)
_ws_wake_event.clear()
def run_alarm_ws_consumer_loop(
client: UMEClient,
*,
on_status: Callable[[str], None] | None = None,
stop_event: threading.Event | None = None,
is_paused: Callable[[], bool] | None = None,
) -> None:
"""
Connect WSS only when a persisted/manual subscription exists.
Never auto-establish subscription; reconnect same uri after disconnect.
"""
backoff_s = 2.0
max_backoff_s = 120.0
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")
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")
_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")
time.sleep(2.0)
continue
_run_ws_session(
client,
wss_uri=wss_uri,
subscription_id=subscription_id,
on_status=on_status,
stop_event=stop_event,
)
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")
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]}")
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]}")
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")
continue
if on_status is not None:
on_status(f"reconnect_ws_in_{int(backoff_s)}s")
slept = 0.0
while slept < backoff_s:
if stop_event is not None and stop_event.is_set():
return
if _ws_wake_event.is_set():
_ws_wake_event.clear()
break
time.sleep(0.5)
slept += 0.5
backoff_s = min(max_backoff_s, backoff_s * 2.0)
def shutdown_ws_consumer() -> None:
"""Process exit: close WSS only; subscription remains on UME until user cancels."""
request_ws_reconnect()
_ws_wake_event.set()
def start_ume_alarm_ws_consumer(
client: UMEClient,
*,
on_status: Callable[[str], None] | None = None,
stop_event: threading.Event | None = None,
is_paused: Callable[[], bool] | None = None,
) -> threading.Thread:
ev = stop_event if stop_event is not None else threading.Event()
def _target() -> None:
if not bool(getattr(settings, "ume_alarm_ws_enabled", True)):
return
if not str(client.base_url or "").strip():
if on_status is not None:
on_status("disabled:no_base_url")
return
global _active_client
with _shutdown_lock:
_active_client = client
run_alarm_ws_consumer_loop(
client,
on_status=on_status,
stop_event=ev,
is_paused=is_paused,
)
thread = threading.Thread(target=_target, name="ume-alarm-ws-consumer", daemon=True)
thread.start()
return thread

View file

@ -53,6 +53,9 @@ class UMEClient:
token_logout_path: str | None = None, token_logout_path: str | None = None,
ne_path: str | None = None, ne_path: str | None = None,
alarms_path: str | None = None, alarms_path: str | None = None,
notification_establish_path: str | None = None,
notification_delete_path: str | None = None,
notification_topic: str | None = None,
token_loader: Callable[[], tuple[str, float] | None] | None = None, token_loader: Callable[[], tuple[str, float] | None] | None = None,
token_saver: Callable[[str, float], None] | None = None, token_saver: Callable[[str, float], None] | None = None,
token_clearer: Callable[[], None] | None = None, token_clearer: Callable[[], None] | None = None,
@ -78,6 +81,19 @@ class UMEClient:
self.token_logout_path = str(token_logout_path if token_logout_path is not None else settings.ume_token_logout_path).strip() self.token_logout_path = str(token_logout_path if token_logout_path is not None else settings.ume_token_logout_path).strip()
self.ne_path = str(ne_path if ne_path is not None else settings.ume_ne_path).strip() self.ne_path = str(ne_path if ne_path is not None else settings.ume_ne_path).strip()
self.alarms_path = str(alarms_path if alarms_path is not None else settings.ume_alarms_path).strip() self.alarms_path = str(alarms_path if alarms_path is not None else settings.ume_alarms_path).strip()
self.notification_establish_path = str(
notification_establish_path
if notification_establish_path is not None
else settings.ume_notification_establish_path
).strip()
self.notification_delete_path = str(
notification_delete_path
if notification_delete_path is not None
else settings.ume_notification_delete_path
).strip()
self.notification_topic = str(
notification_topic if notification_topic is not None else settings.ume_notification_topic
).strip() or "ALARM"
self._token_loader = token_loader self._token_loader = token_loader
self._token_saver = token_saver self._token_saver = token_saver
self._token_clearer = token_clearer self._token_clearer = token_clearer
@ -373,6 +389,56 @@ class UMEClient:
return self.renew_token() return self.renew_token()
return self.login(force=False) return self.login(force=False)
def _request_json_with_current_token(
self,
method: str,
path: str,
*,
params: dict[str, Any] | None = None,
body: dict[str, Any] | None = None,
) -> tuple[dict[str, Any], RequestDiagnostics]:
"""REST call with Content-Type + accessToken from memory/store only (no login/renew)."""
if not self.has_valid_token():
raise RuntimeError("ume_no_valid_token")
url = self._build_url(path)
m = str(method or "GET").upper()
t0 = time()
try:
with self._client() as client:
resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True))
marker = str(resp.headers.get("marker") or "").strip()
is_end_raw = str(resp.headers.get("is-end-of-reply") or "").strip().lower()
is_end_of_reply: bool | None = None
if is_end_raw in {"true", "false"}:
is_end_of_reply = is_end_raw == "true"
if not resp.is_success:
diag = RequestDiagnostics(
method=m,
path=path,
status_code=int(resp.status_code),
latency_ms=int((time() - t0) * 1000),
retry_count=0,
error_code=f"http_{int(resp.status_code)}",
marker=marker,
is_end_of_reply=is_end_of_reply,
)
raise RuntimeError(f"ume_request_failed:{resp.status_code}:{resp.text[:240]}")
data = _coerce_dict(resp.json())
diag = RequestDiagnostics(
method=m,
path=path,
status_code=int(resp.status_code),
latency_ms=int((time() - t0) * 1000),
retry_count=0,
marker=marker,
is_end_of_reply=is_end_of_reply,
)
return data, diag
except Exception as exc:
if isinstance(exc, RuntimeError):
raise
raise RuntimeError(f"ume_request_failed:{str(exc)[:240]}") from exc
def request_json( def request_json(
self, self,
method: str, method: str,
@ -520,3 +586,68 @@ class UMEClient:
data, diag = self.request_json("GET", self.alarms_path, params=params) data, diag = self.request_json("GET", self.alarms_path, params=params)
rows = self._extract_named_list(data, ["alarm-list", "alarm"]) rows = self._extract_named_list(data, ["alarm-list", "alarm"])
return rows, diag return rows, diag
def _extract_subscription_output(self, payload: dict[str, Any]) -> tuple[str, str]:
sub_id = ""
uri = ""
def walk(node: Any) -> None:
nonlocal sub_id, uri
if isinstance(node, dict):
for k, v in node.items():
key = str(k).lower()
if not sub_id and key == "id" and isinstance(v, str):
sub_id = v.strip()
if not uri and key == "uri" and isinstance(v, str):
uri = v.strip()
walk(v)
elif isinstance(node, list):
for item in node:
walk(item)
walk(payload)
return sub_id, uri
def establish_alarm_subscription(self, *, topic: str | None = None) -> tuple[str, str]:
"""POST establish-subscription using token from shared store (no login/renew in this path)."""
self._sync_token_from_store()
if not self.has_valid_token():
raise RuntimeError("ume_establish_subscription_no_valid_token")
topic_val = str(topic if topic is not None else self.notification_topic).strip() or "ALARM"
body = {"input": {"topic": topic_val}}
data, _diag = self._request_json_with_current_token("POST", self.notification_establish_path, body=body)
sub_id, uri = self._extract_subscription_output(data)
if not sub_id or not uri:
raise RuntimeError("ume_establish_subscription_failed:missing_id_or_uri")
return sub_id, uri
def delete_alarm_subscription(self, subscription_id: str) -> None:
sub_id = str(subscription_id or "").strip()
if not sub_id:
return
self._sync_token_from_store()
if not self.has_valid_token():
return
body = {"input": {"id": sub_id}}
try:
self._request_json_with_current_token("POST", self.notification_delete_path, body=body)
except Exception:
return
def has_valid_token(self) -> bool:
"""True when a non-expired token is present in memory (call _sync_token_from_store first)."""
token = self._token_value.strip()
if not token:
return False
now = time()
return now < (self._token_expires_at - self.token_refresh_skew_s)
def ws_auth_headers(self) -> dict[str, str]:
"""
Headers for WSS handshake — same fields as REST (_headers).
Does not login/renew; relies on shared token store (token_keepalive / other sync paths).
"""
self._sync_token_from_store()
if not self.has_valid_token():
raise RuntimeError("ume_ws_no_valid_token")
return self._headers(include_token=True)

View file

@ -166,6 +166,106 @@ def _derive_ne_id_from_alarm(alarm: dict[str, Any]) -> str:
return "" return ""
def _normalize_yang_key(key: str) -> str:
raw = str(key or "").strip()
if not raw:
return ""
if ":" in raw:
return raw.rsplit(":", 1)[-1]
return raw
def normalize_yang_alarm(raw: dict[str, Any]) -> dict[str, Any]:
"""Flatten YANG namespace-prefixed keys (e.g. zte-alarms:alarmkey) for _pick()."""
out: dict[str, Any] = {}
for k, v in raw.items():
nk = _normalize_yang_key(str(k))
if not nk:
continue
if nk in out and out[nk] not in (None, ""):
continue
out[nk] = v
return out
def _is_alarm_cleared(alarm: dict[str, Any]) -> bool:
val = _pick(alarm, "isCleared", "is-cleared")
if isinstance(val, bool):
return val
text = _s(val).lower()
return text in {"true", "1", "yes"}
def extract_alarm_from_notification(payload: dict[str, Any]) -> dict[str, Any] | None:
"""Parse alarm-notification from a WS/REST notification envelope."""
if not isinstance(payload, dict):
return None
def _find_alarm_notification(node: Any) -> dict[str, Any] | None:
if isinstance(node, dict):
for k, v in node.items():
key = str(k).lower()
if key in {"alarm-notification", "alarm_notification"} and isinstance(v, dict):
return normalize_yang_alarm(v)
found = _find_alarm_notification(v)
if found is not None:
return found
elif isinstance(node, list):
for item in node:
found = _find_alarm_notification(item)
if found is not None:
return found
return None
direct = _find_alarm_notification(payload)
if direct is not None:
return direct
return normalize_yang_alarm(payload) if payload else None
def apply_alarm_to_current(db: Session, alarm: dict[str, Any], *, touch_ts: datetime) -> tuple[str, bool]:
"""
Apply one alarm to ume_alarms_current.
Returns (action, changed) where action is inserted|updated|deleted|skipped.
"""
norm = normalize_yang_alarm(alarm) if alarm else {}
if not norm:
return "skipped", False
if _is_alarm_cleared(norm):
key = _alarm_key(norm)
if not key:
return "skipped", False
existing = db.get(UmeAlarmCurrent, key)
if existing is None:
return "deleted", False
db.delete(existing)
return "deleted", True
key = _alarm_key(norm)
existing = db.get(UmeAlarmCurrent, key)
if existing is None:
existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=touch_ts)
db.add(existing)
action = "inserted"
else:
action = "updated"
existing.ne_id = _s(_derive_ne_id_from_alarm(norm))
existing.host_name = _lookup_host_name(db, existing.ne_id)
existing.object_name = _s(_pick(norm, "objectName", "object-name"))
existing.event_type = _s(_pick(norm, "eventType", "event-type"))
existing.native_probable_cause = _s(_pick(norm, "nativeProbableCause", "native-probable-cause"))
existing.perceived_severity = _s(_pick(norm, "perceivedSeverity", "perceived-severity"))
existing.is_cleared = _s(_pick(norm, "isCleared", "is-cleared"))
existing.time_created = _s(_pick(norm, "timeCreated", "time-created"))
existing.root_cause_alarm_indication = _s(
_pick(norm, "rootCauseAlarmIndication", "root-cause-alarm-indication")
)
existing.last_seen_at = touch_ts
existing.raw_json = json.dumps(norm, ensure_ascii=False, default=str)
return action, True
def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob:
return UmeSyncJob( return UmeSyncJob(
domain=domain, domain=domain,
@ -414,10 +514,9 @@ def _sync_alarms_common(
max_pages = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000) max_pages = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000)
max_pages = max(1, min(max_pages, 20000)) max_pages = max(1, min(max_pages, 20000))
def upsert_alarm(alarm: dict[str, Any], *, touch_ts: datetime) -> None: def upsert_alarm_history(alarm: dict[str, Any], *, touch_ts: datetime) -> None:
nonlocal inserted, updated nonlocal inserted, updated
key = _alarm_key(alarm) key = _alarm_key(alarm)
if is_uncleared:
existing = db.get(UmeAlarmHistory, key) existing = db.get(UmeAlarmHistory, key)
if existing is None: if existing is None:
existing = UmeAlarmHistory(alarm_key=key, first_seen_at=touch_ts) existing = UmeAlarmHistory(alarm_key=key, first_seen_at=touch_ts)
@ -425,14 +524,6 @@ def _sync_alarms_common(
inserted += 1 inserted += 1
else: else:
updated += 1 updated += 1
else:
existing = db.get(UmeAlarmCurrent, key)
if existing is None:
existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=touch_ts)
db.add(existing)
inserted += 1
else:
updated += 1
existing.ne_id = _s(_derive_ne_id_from_alarm(alarm)) existing.ne_id = _s(_derive_ne_id_from_alarm(alarm))
existing.host_name = _lookup_host_name(db, existing.ne_id) existing.host_name = _lookup_host_name(db, existing.ne_id)
existing.object_name = _s(_pick(alarm, "objectName", "object-name")) existing.object_name = _s(_pick(alarm, "objectName", "object-name"))
@ -457,7 +548,14 @@ def _sync_alarms_common(
for rows in pages: for rows in pages:
pulled += len(rows) pulled += len(rows)
for alarm in rows: for alarm in rows:
upsert_alarm(alarm, touch_ts=sync_batch_ts) if is_uncleared:
upsert_alarm_history(alarm, touch_ts=sync_batch_ts)
else:
action, _changed = apply_alarm_to_current(db, alarm, touch_ts=sync_batch_ts)
if action == "inserted":
inserted += 1
elif action == "updated":
updated += 1
db.flush() db.flush()
page_no = int(meta.get("page_count") or 0) page_no = int(meta.get("page_count") or 0)
next_marker = str(meta.get("last_marker") or "") next_marker = str(meta.get("last_marker") or "")

View file

@ -9,3 +9,4 @@ PyYAML>=6.0.0
pydantic>=2.8.0 pydantic>=2.8.0
pydantic-settings>=2.3.0 pydantic-settings>=2.3.0
python-multipart>=0.0.9 python-multipart>=0.0.9
websocket-client>=1.8.0

View file

@ -1,6 +1,7 @@
from __future__ import annotations from __future__ import annotations
import unittest import unittest
from typing import Any
from unittest.mock import patch from unittest.mock import patch
from sqlalchemy import create_engine from sqlalchemy import create_engine
@ -10,7 +11,24 @@ from netx_api.db import Base
from netx_api.main import _extract_ume_raw_group_field, _serialize_ume_alarm_raw_row, sql_ume_query, ume_alarms_fields from netx_api.main import _extract_ume_raw_group_field, _serialize_ume_alarm_raw_row, sql_ume_query, ume_alarms_fields
from netx_api.models import UmeAlarmCurrent, UmeInventoryNE from netx_api.models import UmeAlarmCurrent, UmeInventoryNE
from netx_api.ume_client import UMEClient from netx_api.ume_client import UMEClient
from netx_api.ume_sync_service import _derive_ne_id_from_alarm, sync_alarms_current, sync_inventory_full 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 (
cancel_alarm_subscription_manual,
establish_alarm_subscription_manual,
get_subscription_status,
load_persisted_subscription,
process_alarm_notification,
)
from netx_api.ume_sync_service import (
_derive_ne_id_from_alarm,
apply_alarm_to_current,
extract_alarm_from_notification,
normalize_yang_alarm,
sync_alarms_current,
sync_inventory_full,
)
from fastapi import HTTPException from fastapi import HTTPException
@ -127,6 +145,53 @@ class UMEClientTests(unittest.TestCase):
self.assertEqual(len(rows), 2) self.assertEqual(len(rows), 2)
self.assertEqual(rows[1].get("alarmKey"), "AK-2") self.assertEqual(rows[1].get("alarmKey"), "AK-2")
def test_establish_alarm_subscription(self):
from time import time as _time
client = UMEClient(base_url="https://ume.local:18014", username="u", password="p", verify_tls=False)
client._token_value = "token-1"
client._token_expires_at = _time() + 3600
seen: dict[str, Any] = {}
def _fake_request(method: str, path: str, *, params=None, body=None):
seen["method"] = method
seen["path"] = path
seen["body"] = body
return (
{
"output": {
"id": "3282ac78-b38a-4242-81d2-cc5b77c28ef8",
"uri": "wss://ume.local:18014/restconf/stream/3282ac78-b38a-4242-81d2-cc5b77c28ef8",
}
},
None,
)
client._request_json_with_current_token = _fake_request # type: ignore[method-assign]
sub_id, uri = client.establish_alarm_subscription()
self.assertEqual(seen["method"], "POST")
self.assertEqual(seen["body"], {"input": {"topic": "ALARM"}})
self.assertEqual(sub_id, "3282ac78-b38a-4242-81d2-cc5b77c28ef8")
self.assertIn("wss://", uri)
def test_delete_alarm_subscription(self):
from time import time as _time
client = UMEClient(base_url="https://ume.local:18014", username="u", password="p", verify_tls=False)
client._token_value = "token-1"
client._token_expires_at = _time() + 3600
seen: dict[str, Any] = {}
def _fake_request(method: str, path: str, *, params=None, body=None):
seen["method"] = method
seen["body"] = body
return ({}, None)
client._request_json_with_current_token = _fake_request # type: ignore[method-assign]
client.delete_alarm_subscription("sub-to-delete")
self.assertEqual(seen["method"], "POST")
self.assertEqual(seen["body"], {"input": {"id": "sub-to-delete"}})
class UmeSyncServiceTests(unittest.TestCase): class UmeSyncServiceTests(unittest.TestCase):
def setUp(self): def setUp(self):
@ -570,5 +635,129 @@ class UmeSyncServiceTests(unittest.TestCase):
self.assertIn("alarm_host_name", set(data["selectable_fields"])) self.assertIn("alarm_host_name", set(data["selectable_fields"]))
class UmeAlarmNotificationTests(unittest.TestCase):
def setUp(self):
engine = create_engine("sqlite+pysqlite:///:memory:", future=True)
TestingSessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False, expire_on_commit=False)
Base.metadata.create_all(bind=engine)
self.db = TestingSessionLocal()
def tearDown(self):
self.db.close()
def test_normalize_yang_alarm_strips_prefix(self):
raw = {
"zte-alarms:alarmkey": "AK-WS-1",
"zte-alarms:is-cleared": False,
"zte-alarms:perceivedSeverity": "critical",
}
norm = normalize_yang_alarm(raw)
self.assertEqual(norm.get("alarmkey"), "AK-WS-1")
self.assertEqual(norm.get("is-cleared"), False)
self.assertEqual(norm.get("perceivedSeverity"), "critical")
def test_apply_alarm_to_current_insert_from_notification(self):
payload = {
"alarm-notification": {
"zte-alarms:alarmkey": "AK-WS-2",
"zte-alarms:is-cleared": False,
"zte-alarms:perceivedSeverity": "major",
"zte-alarms:objectName": "ME{00ceb960-1b62-478e-8303-0935ffea1d28}",
"zte-alarms:time-created": "2025-01-03T07:55:00.823+08:00",
}
}
alarm = extract_alarm_from_notification(payload)
self.assertIsNotNone(alarm)
from datetime import datetime
action, changed = apply_alarm_to_current(self.db, alarm or {}, touch_ts=datetime.utcnow())
self.db.commit()
self.assertEqual(action, "inserted")
self.assertTrue(changed)
row = self.db.get(UmeAlarmCurrent, "AK-WS-2")
self.assertIsNotNone(row)
self.assertEqual(row.perceived_severity, "major")
def test_apply_alarm_cleared_deletes_current_only(self):
from datetime import datetime
touch = datetime.utcnow()
self.db.add(
UmeAlarmCurrent(
alarm_key="AK-CLEARED-1",
first_seen_at=touch,
last_seen_at=touch,
is_cleared="false",
)
)
self.db.commit()
alarm = {
"alarmkey": "AK-CLEARED-1",
"is-cleared": True,
}
action, changed = apply_alarm_to_current(self.db, alarm, touch_ts=touch)
self.db.commit()
self.assertEqual(action, "deleted")
self.assertTrue(changed)
self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-1"))
def test_subscription_store_and_manual_establish(self):
from time import time as _time
ume_alarm_ws._clear_active_subscription()
engine = create_engine("sqlite+pysqlite:///:memory:", future=True)
TestingSessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False, expire_on_commit=False)
Base.metadata.create_all(bind=engine)
db = TestingSessionLocal()
client = UMEClient(base_url="https://ume.local:18014", username="u", password="p", verify_tls=False)
client._token_value = "token-1"
client._token_expires_at = _time() + 3600
calls = {"n": 0}
def _fake_establish(*, topic=None):
calls["n"] += 1
return ("sub-1", "wss://ume.local:18014/restconf/stream/sub-1")
client.establish_alarm_subscription = _fake_establish # type: ignore[method-assign]
client.delete_alarm_subscription = lambda _id: None # type: ignore[method-assign]
st = establish_alarm_subscription_manual(client, db)
self.assertTrue(st["active"])
self.assertEqual(st["subscription_id"], "sub-1")
self.assertFalse(st.get("already_exists"))
self.assertEqual(calls["n"], 1)
st2 = establish_alarm_subscription_manual(client, db)
self.assertTrue(st2.get("already_exists"))
self.assertEqual(st2["subscription_id"], "sub-1")
self.assertEqual(calls["n"], 1)
loaded = load_subscription(db)
self.assertIsNotNone(loaded)
self.assertEqual(loaded[0], "sub-1")
self.assertTrue(get_subscription_status()["active"])
cancel_alarm_subscription_manual(client, db)
self.assertFalse(get_subscription_status()["active"])
self.assertIsNone(load_subscription(db))
db.close()
def test_process_alarm_notification_via_ws_helper(self):
from datetime import datetime
payload = {
"alarm-notification": {
"zte-alarms:alarmkey": "AK-WS-3",
"zte-alarms:is-cleared": False,
"zte-alarms:perceivedSeverity": "warning",
}
}
action, changed = process_alarm_notification(self.db, payload)
self.db.commit()
self.assertEqual(action, "inserted")
self.assertTrue(changed)
self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-WS-3"))
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()

View file

@ -5,6 +5,9 @@ import {
disconnectUmeToken, disconnectUmeToken,
fetchUmeCurrentAlarms, fetchUmeCurrentAlarms,
fetchUmeNe, fetchUmeNe,
cancelUmeAlarmSubscription,
establishUmeAlarmSubscription,
fetchUmeAlarmSubscriptionStatus,
fetchUmeSyncStatus, fetchUmeSyncStatus,
fetchUmeTokenStatus, fetchUmeTokenStatus,
refreshUmeToken, refreshUmeToken,
@ -21,6 +24,7 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
const queryClient = useQueryClient(); const queryClient = useQueryClient();
const [tokenOpError, setTokenOpError] = useState(""); const [tokenOpError, setTokenOpError] = useState("");
const [runtimeTaskError, setRuntimeTaskError] = useState(""); const [runtimeTaskError, setRuntimeTaskError] = useState("");
const [subscriptionOpError, setSubscriptionOpError] = useState("");
const [syncPage, setSyncPage] = useState(1); const [syncPage, setSyncPage] = useState(1);
const [syncPageSize, setSyncPageSize] = useState(20); const [syncPageSize, setSyncPageSize] = useState(20);
const [neKeyword, setNeKeyword] = useState(""); const [neKeyword, setNeKeyword] = useState("");
@ -79,6 +83,53 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
toastError?.(msg); toastError?.(msg);
}, },
}); });
const subscriptionStatusQuery = useQuery({
queryKey: ["umeAlarmSubscription"],
queryFn: fetchUmeAlarmSubscriptionStatus,
staleTime: 3000,
refetchInterval: 5000,
});
const subscriptionEstablishMutation = useMutation({
mutationFn: establishUmeAlarmSubscription,
onMutate: () => setSubscriptionOpError(""),
onSuccess: async (res) => {
if (!res?.active) {
const msg = "建立订阅未返回有效 id/uri";
setSubscriptionOpError(msg);
toastError?.(msg);
} else {
setSubscriptionOpError("");
toastOk?.(
res.already_exists
? "订阅已存在,未重复建立(WSS 将保持/恢复连接)"
: "告警订阅已建立,WSS 将自动连接",
);
}
await queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] });
await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] });
},
onError: (err) => {
const msg = String(err);
setSubscriptionOpError(msg);
toastError?.(msg);
},
});
const subscriptionCancelMutation = useMutation({
mutationFn: cancelUmeAlarmSubscription,
onMutate: () => setSubscriptionOpError(""),
onSuccess: async () => {
setSubscriptionOpError("");
toastOk?.("告警订阅已取消");
await queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] });
await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] });
},
onError: (err) => {
const msg = String(err);
setSubscriptionOpError(msg);
toastError?.(msg);
},
});
const tokenDisconnectMutation = useMutation({ const tokenDisconnectMutation = useMutation({
mutationFn: disconnectUmeToken, mutationFn: disconnectUmeToken,
onMutate: () => { onMutate: () => {
@ -135,6 +186,14 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
}); });
const runningTasks = (syncStatusQuery.data?.items || []).filter((x) => String(x.status || "").toLowerCase() === "running"); const runningTasks = (syncStatusQuery.data?.items || []).filter((x) => String(x.status || "").toLowerCase() === "running");
const runtimeTasks = syncStatusQuery.data?.runtime_tasks || []; const runtimeTasks = syncStatusQuery.data?.runtime_tasks || [];
const alarmSub =
subscriptionStatusQuery.data ??
syncStatusQuery.data?.alarm_subscription ??
({ active: false } as const);
const wsConsumer = runtimeTasks.find((t) => t.task === "alarms_current_ws_consumer");
const subscriptionActive = Boolean(alarmSub.active);
const subPending =
subscriptionEstablishMutation.isPending || subscriptionCancelMutation.isPending;
const runtimeTaskMutation = useMutation({ const runtimeTaskMutation = useMutation({
mutationFn: async (vars: { task: string; action: "pause" | "resume" }) => mutationFn: async (vars: { task: string; action: "pause" | "resume" }) =>
@ -228,6 +287,78 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
</div> </div>
)} )}
</article> </article>
<article className="card card--full">
<h3>UME 告警订阅(WebSocket)</h3>
<p className="muted" style={{ marginTop: 0 }}>
订阅需手动建立/取消;建立后(或重启后若库中仍有有效订阅)后台会自动连接 WSS 接收实时告警。
</p>
<div className="actions-row actions-row--inline">
<span className={`conn-pill conn-pill--${subscriptionActive ? "up" : "down"}`}>
订阅: {subscriptionActive ? "已建立" : "未建立"}
</span>
{subscriptionActive && alarmSub.subscription_id ? (
<span className="conn-pill" title={alarmSub.wss_uri || ""}>
id: {String(alarmSub.subscription_id).slice(0, 12)}…
</span>
) : null}
{wsConsumer ? (
<span
className={`conn-pill conn-pill--${
String(wsConsumer.last_error || "").includes("connected") ||
wsConsumer.status === "running"
? "up"
: String(wsConsumer.last_error || "").includes("no_subscription")
? "down"
: "unknown"
}`}
title={wsConsumer.last_error || ""}
>
WSS: {wsConsumer.last_error || wsConsumer.status || "-"}
</span>
) : null}
</div>
<div className="actions-row actions-row--inline">
<button
type="button"
onClick={() => subscriptionEstablishMutation.mutate()}
disabled={subPending || subscriptionActive || !hasToken}
title={!hasToken ? "请先登录 UME token" : undefined}
>
{subscriptionEstablishMutation.isPending ? (
<>
<span className="inline-spinner" aria-hidden />
建立订阅…
</>
) : (
"建立告警订阅"
)}
</button>
<button
type="button"
onClick={() => subscriptionCancelMutation.mutate()}
disabled={subPending || !subscriptionActive}
>
{subscriptionCancelMutation.isPending ? (
<>
<span className="inline-spinner" aria-hidden />
取消订阅…
</>
) : (
"取消告警订阅"
)}
</button>
<button
type="button"
onClick={() => queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] })}
disabled={subscriptionStatusQuery.isFetching}
>
刷新订阅状态
</button>
</div>
{subscriptionOpError ? (
<div className="pill pill--high">订阅操作失败: {subscriptionOpError}</div>
) : null}
</article>
<article className="card card--full"> <article className="card card--full">
<h3>UME 同步</h3> <h3>UME 同步</h3>
<div className="actions-row actions-row--inline"> <div className="actions-row actions-row--inline">

View file

@ -6,6 +6,7 @@ import type {
ImportHistoryItem, ImportHistoryItem,
IntegrationStatus, IntegrationStatus,
UmeAlarmItem, UmeAlarmItem,
UmeAlarmSubscriptionStatus,
UmeNeItem, UmeNeItem,
UmeSyncStatusResponse, UmeSyncStatusResponse,
UmeTokenStatus, UmeTokenStatus,
@ -100,6 +101,15 @@ export const fetchAlarms = (params: {
return apiGet<AlarmQueryResponse>(`/v1/alarms?${p.toString()}`); return apiGet<AlarmQueryResponse>(`/v1/alarms?${p.toString()}`);
}; };
export const fetchUmeAlarmSubscriptionStatus = () =>
apiGet<UmeAlarmSubscriptionStatus>("/v1/ume/alarm-subscription/status");
export const establishUmeAlarmSubscription = () =>
apiPost<UmeAlarmSubscriptionStatus>("/v1/ume/alarm-subscription/establish", {});
export const cancelUmeAlarmSubscription = () =>
apiPost<UmeAlarmSubscriptionStatus>("/v1/ume/alarm-subscription/cancel", {});
export const fetchUmeSyncStatus = (params: { page: number; pageSize: number }) => { export const fetchUmeSyncStatus = (params: { page: number; pageSize: number }) => {
const p = new URLSearchParams(); const p = new URLSearchParams();
p.set("page", String(Math.max(1, Number(params.page || 1)))); p.set("page", String(Math.max(1, Number(params.page || 1))));

View file

@ -89,12 +89,26 @@ export type UmeSyncJobItem = {
ended_at?: string | null; ended_at?: string | null;
}; };
export type UmeAlarmSubscriptionStatus = {
ok?: boolean;
created?: boolean;
already_exists?: boolean;
active: boolean;
subscription_id?: string;
wss_uri?: string;
topic?: string;
ws_consumer_status?: string;
ws_consumer_last_error?: string;
ws_consumer_last_run_at?: string | null;
};
export type UmeSyncStatusResponse = { export type UmeSyncStatusResponse = {
total?: number; total?: number;
page?: number; page?: number;
page_size?: number; page_size?: number;
items: UmeSyncJobItem[]; items: UmeSyncJobItem[];
latest_by_domain?: Record<string, UmeSyncJobItem>; latest_by_domain?: Record<string, UmeSyncJobItem>;
alarm_subscription?: UmeAlarmSubscriptionStatus;
runtime_tasks?: Array<{ runtime_tasks?: Array<{
task: string; task: string;
status: string; status: string;