diff --git a/.env.example b/.env.example index 07bbc05..8faa188 100644 --- a/.env.example +++ b/.env.example @@ -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_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list 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 diff --git a/netx_api/config.py b/netx_api/config.py index 58da2fb..d042d7e 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -41,7 +41,11 @@ class Settings(BaseSettings): ume_keepalive_interval_s: int = 600 ume_keepalive_renew_before_s: int = 900 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_every_hours: int = 48 ume_token_path: str = "/restconf/operations/zte-security:oauth_token" diff --git a/netx_api/main.py b/netx_api/main.py index de4562a..c819d56 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -34,6 +34,14 @@ from .models import ( from .models import ImportJob from .parser_config import load_parser_config 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_token_store import ( clear_shared_token, @@ -73,8 +81,10 @@ _SQL_FORBIDDEN_RE = re.compile( _UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = { "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_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": ""}, } +_UME_WS_STOP_EVENT: threading.Event | None = None _UME_RUNTIME_PAUSED: dict[str, bool] = {} UME_KNOWN_RUNTIME_TASKS: tuple[str, ...] = tuple(_UME_RUNTIME_TASKS.keys()) _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 not bool(getattr(settings, "ume_sync_alarms_current_enabled", True)): 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)) 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 not bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)): return None, "未启用" @@ -701,7 +715,7 @@ def on_startup() -> None: ) try: 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)) def _alarms_current_sync_loop() -> None: @@ -832,6 +846,46 @@ def on_startup() -> None: last_run_at=datetime.now(timezone.utc), 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") @@ -872,6 +926,41 @@ def ume_token_disconnect() -> dict[str, Any]: 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") def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict[str, Any]: body = payload or {} @@ -993,6 +1082,7 @@ def ume_sync_status( "items": items, "latest_by_domain": latest_by_domain, "runtime_tasks": _list_runtime_tasks(), + "alarm_subscription": get_subscription_status(), } diff --git a/netx_api/models.py b/netx_api/models.py index aeb50ce..b6a8e70 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -199,6 +199,19 @@ class UmeAlarmHistory(Base): 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): __tablename__ = "ume_token_cache" diff --git a/netx_api/ume_alarm_subscription_store.py b/netx_api/ume_alarm_subscription_store.py new file mode 100644 index 0000000..411127b --- /dev/null +++ b/netx_api/ume_alarm_subscription_store.py @@ -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() diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py new file mode 100644 index 0000000..224a911 --- /dev/null +++ b/netx_api/ume_alarm_ws.py @@ -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 diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index 0c6dac6..3825f50 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -53,6 +53,9 @@ class UMEClient: token_logout_path: str | None = None, ne_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_saver: Callable[[str, float], 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.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.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_saver = token_saver self._token_clearer = token_clearer @@ -373,6 +389,56 @@ class UMEClient: return self.renew_token() 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( self, method: str, @@ -520,3 +586,68 @@ class UMEClient: data, diag = self.request_json("GET", self.alarms_path, params=params) rows = self._extract_named_list(data, ["alarm-list", "alarm"]) 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) diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 56ce50f..b469b25 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -166,6 +166,106 @@ def _derive_ne_id_from_alarm(alarm: dict[str, Any]) -> str: 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: return UmeSyncJob( domain=domain, @@ -414,25 +514,16 @@ def _sync_alarms_common( 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)) - 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 key = _alarm_key(alarm) - if is_uncleared: - existing = db.get(UmeAlarmHistory, key) - if existing is None: - existing = UmeAlarmHistory(alarm_key=key, first_seen_at=touch_ts) - db.add(existing) - inserted += 1 - else: - updated += 1 + existing = db.get(UmeAlarmHistory, key) + if existing is None: + existing = UmeAlarmHistory(alarm_key=key, first_seen_at=touch_ts) + db.add(existing) + inserted += 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 + updated += 1 existing.ne_id = _s(_derive_ne_id_from_alarm(alarm)) existing.host_name = _lookup_host_name(db, existing.ne_id) existing.object_name = _s(_pick(alarm, "objectName", "object-name")) @@ -457,7 +548,14 @@ def _sync_alarms_common( for rows in pages: pulled += len(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() page_no = int(meta.get("page_count") or 0) next_marker = str(meta.get("last_marker") or "") diff --git a/requirements.txt b/requirements.txt index 8265de0..a0917ba 100644 --- a/requirements.txt +++ b/requirements.txt @@ -9,3 +9,4 @@ PyYAML>=6.0.0 pydantic>=2.8.0 pydantic-settings>=2.3.0 python-multipart>=0.0.9 +websocket-client>=1.8.0 diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 5093339..3c689b2 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -1,6 +1,7 @@ from __future__ import annotations import unittest +from typing import Any from unittest.mock import patch 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.models import UmeAlarmCurrent, UmeInventoryNE 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 @@ -127,6 +145,53 @@ class UMEClientTests(unittest.TestCase): self.assertEqual(len(rows), 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): def setUp(self): @@ -570,5 +635,129 @@ class UmeSyncServiceTests(unittest.TestCase): 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__": unittest.main() diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index 6ac0fa0..52f71b7 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -5,6 +5,9 @@ import { disconnectUmeToken, fetchUmeCurrentAlarms, fetchUmeNe, + cancelUmeAlarmSubscription, + establishUmeAlarmSubscription, + fetchUmeAlarmSubscriptionStatus, fetchUmeSyncStatus, fetchUmeTokenStatus, refreshUmeToken, @@ -21,6 +24,7 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { const queryClient = useQueryClient(); const [tokenOpError, setTokenOpError] = useState(""); const [runtimeTaskError, setRuntimeTaskError] = useState(""); + const [subscriptionOpError, setSubscriptionOpError] = useState(""); const [syncPage, setSyncPage] = useState(1); const [syncPageSize, setSyncPageSize] = useState(20); const [neKeyword, setNeKeyword] = useState(""); @@ -79,6 +83,53 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { 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({ mutationFn: disconnectUmeToken, onMutate: () => { @@ -135,6 +186,14 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { }); const runningTasks = (syncStatusQuery.data?.items || []).filter((x) => String(x.status || "").toLowerCase() === "running"); 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({ mutationFn: async (vars: { task: string; action: "pause" | "resume" }) => @@ -228,6 +287,78 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { )} +
+

UME 告警订阅(WebSocket)

+

+ 订阅需手动建立/取消;建立后(或重启后若库中仍有有效订阅)后台会自动连接 WSS 接收实时告警。 +

+
+ + 订阅: {subscriptionActive ? "已建立" : "未建立"} + + {subscriptionActive && alarmSub.subscription_id ? ( + + id: {String(alarmSub.subscription_id).slice(0, 12)}… + + ) : null} + {wsConsumer ? ( + + WSS: {wsConsumer.last_error || wsConsumer.status || "-"} + + ) : null} +
+
+ + + +
+ {subscriptionOpError ? ( +
订阅操作失败: {subscriptionOpError}
+ ) : null} +

UME 同步

diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 673ccc1..63b5d88 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -6,6 +6,7 @@ import type { ImportHistoryItem, IntegrationStatus, UmeAlarmItem, + UmeAlarmSubscriptionStatus, UmeNeItem, UmeSyncStatusResponse, UmeTokenStatus, @@ -100,6 +101,15 @@ export const fetchAlarms = (params: { return apiGet(`/v1/alarms?${p.toString()}`); }; +export const fetchUmeAlarmSubscriptionStatus = () => + apiGet("/v1/ume/alarm-subscription/status"); + +export const establishUmeAlarmSubscription = () => + apiPost("/v1/ume/alarm-subscription/establish", {}); + +export const cancelUmeAlarmSubscription = () => + apiPost("/v1/ume/alarm-subscription/cancel", {}); + export const fetchUmeSyncStatus = (params: { page: number; pageSize: number }) => { const p = new URLSearchParams(); p.set("page", String(Math.max(1, Number(params.page || 1)))); diff --git a/web/src/types.ts b/web/src/types.ts index ef01a10..dc8db9d 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -89,12 +89,26 @@ export type UmeSyncJobItem = { 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 = { total?: number; page?: number; page_size?: number; items: UmeSyncJobItem[]; latest_by_domain?: Record; + alarm_subscription?: UmeAlarmSubscriptionStatus; runtime_tasks?: Array<{ task: string; status: string;