diff --git a/netx_api/main.py b/netx_api/main.py index 00ec624..0d0567a 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -36,6 +36,7 @@ from .parser_config import load_parser_config from .ume_client import UMEClient from .ume_alarm_ws import ( cancel_alarm_subscription_manual, + clear_local_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, get_ws_connection_status, @@ -946,10 +947,15 @@ def ume_alarm_subscription_status(limit: int = 80) -> dict[str, Any]: @app.post("/v1/ume/alarm-subscription/establish") -def ume_alarm_subscription_establish(db: Session = Depends(get_db)) -> dict[str, Any]: +def ume_alarm_subscription_establish( + payload: dict[str, Any] | None = None, + db: Session = Depends(get_db), +) -> dict[str, Any]: client = _ume_client() + body = payload or {} + force_reestablish = bool(body.get("force_reestablish")) try: - st = establish_alarm_subscription_manual(client, db) + st = establish_alarm_subscription_manual(client, db, force_reestablish=force_reestablish) return {"ok": True, "created": not bool(st.get("already_exists")), **st} except Exception as exc: msg = str(exc)[:240] @@ -957,16 +963,33 @@ def ume_alarm_subscription_establish(db: Session = Depends(get_db)) -> dict[str, @app.post("/v1/ume/alarm-subscription/cancel") -def ume_alarm_subscription_cancel(db: Session = Depends(get_db)) -> dict[str, Any]: +def ume_alarm_subscription_cancel( + payload: dict[str, Any] | None = None, + db: Session = Depends(get_db), +) -> dict[str, Any]: client = _ume_client() + body = payload or {} + force_clear_local = bool(body.get("force_clear_local")) try: - st = cancel_alarm_subscription_manual(client, db) + st = cancel_alarm_subscription_manual(client, db, force_clear_local=force_clear_local) + if st.get("needs_local_cleanup"): + return st 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/alarm-subscription/clear-local") +def ume_alarm_subscription_clear_local(db: Session = Depends(get_db)) -> dict[str, Any]: + try: + st = clear_local_alarm_subscription_manual(db) + return {"ok": True, "cleared_local": 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 {} diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py index a7177a4..1877441 100644 --- a/netx_api/ume_alarm_ws.py +++ b/netx_api/ume_alarm_ws.py @@ -44,6 +44,22 @@ _WS_LOG_MAX_RETURN = 100 _ORPHAN_SUB_ID_RE = re.compile(r"id:([0-9a-fA-F-]{8}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{4}-[0-9a-fA-F-]{12})") +_SUBSCRIPTION_MISSING_MARKERS: tuple[str, ...] = ( + "subscription not exist", + "subscription not found", + "not found or overtime", + "please establish again", + "status-code: 404", + "status-reason: subscription not found", + "non-101 status: 503", +) + + +def is_ume_subscription_missing_error(message: str) -> bool: + """True when UME reports the notification subscription is gone or expired.""" + low = str(message or "").lower() + return any(marker in low for marker in _SUBSCRIPTION_MISSING_MARKERS) + def _utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat() @@ -98,6 +114,9 @@ _subscription_topic: str = "ALARM" _active_client: UMEClient | None = None _ws_connection_state: str = "init" _ws_connection_detail: str = "" +_ume_lost_lock = threading.Lock() +_ume_subscription_lost = False +_ume_subscription_lost_reason: str = "" _WS_CONNECTION_LABELS: dict[str, str] = { "init": "初始化", @@ -109,9 +128,42 @@ _WS_CONNECTION_LABELS: dict[str, str] = { "paused": "已暂停", "error": "连接异常", "reconnecting": "重连等待", + "subscription_lost": "UME订阅已丢失", } +def mark_ume_subscription_lost(reason: str) -> None: + global _ume_subscription_lost, _ume_subscription_lost_reason + detail = str(reason or "").strip()[:500] + with _ume_lost_lock: + _ume_subscription_lost = True + _ume_subscription_lost_reason = detail + append_ws_log(f"UME 侧订阅已丢失: {detail[:200]}", level="warning") + request_ws_reconnect() + + +def clear_ume_subscription_lost_flag() -> None: + global _ume_subscription_lost, _ume_subscription_lost_reason + with _ume_lost_lock: + _ume_subscription_lost = False + _ume_subscription_lost_reason = "" + + +def is_ume_subscription_lost() -> bool: + with _ume_lost_lock: + return bool(_ume_subscription_lost) + + +def _server_subscription_lost_fields() -> dict[str, Any]: + with _ume_lost_lock: + lost = bool(_ume_subscription_lost) + reason = str(_ume_subscription_lost_reason or "") + return { + "server_subscription_lost": lost, + "server_subscription_lost_reason": reason, + } + + def _set_ws_connection_state(state: str, *, detail: str = "") -> None: global _ws_connection_state, _ws_connection_detail with _ws_connection_lock: @@ -227,15 +279,35 @@ def get_subscription_status() -> dict[str, Any]: "subscription_id": sub_id, "wss_uri": uri, "topic": topic, + **_server_subscription_lost_fields(), } +def clear_local_alarm_subscription_manual(db: Session) -> dict[str, Any]: + """Drop persisted/in-memory subscription without calling UME delete.""" + clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY) + db.commit() + _clear_active_subscription() + clear_ume_subscription_lost_flag() + request_ws_reconnect() + append_ws_log("cleared local subscription record (UME delete skipped)") + return get_subscription_status() + + def request_ws_reconnect() -> None: _ws_wake_event.set() -def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str, Any]: +def establish_alarm_subscription_manual( + client: UMEClient, + db: Session, + *, + force_reestablish: bool = False, +) -> dict[str, Any]: """Manually establish ALARM subscription (UME API + persist). Does not open WSS.""" + if force_reestablish or is_ume_subscription_lost(): + clear_local_alarm_subscription_manual(db) + existing = load_subscription(db) if existing is not None: sub_id, uri, topic = existing @@ -284,27 +356,62 @@ def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[ with _shutdown_lock: _active_client = client _set_active_subscription(sub_id, uri, topic=topic) + clear_ume_subscription_lost_flag() request_ws_reconnect() append_ws_log(f"manual establish subscription id={sub_id}", subscription_id=sub_id) return {**get_subscription_status(), "already_exists": False} -def cancel_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str, Any]: +def cancel_alarm_subscription_manual( + client: UMEClient, + db: Session, + *, + force_clear_local: bool = False, +) -> 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 "") + ume_already_missing = False if sub_id: client.refresh_if_needed() - _delete_subscription_on_ume(client, sub_id) + try: + _delete_subscription_on_ume(client, sub_id) + except RuntimeError as exc: + if is_ume_subscription_missing_error(str(exc)): + ume_already_missing = True + mark_ume_subscription_lost(str(exc)) + if not force_clear_local: + return { + **get_subscription_status(), + "ok": False, + "ume_already_missing": True, + "needs_local_cleanup": True, + "message": ( + "UME 侧告警订阅已不存在或已过期(服务器订阅已丢失)。" + "请确认是否清除本地订阅记录,然后重新建立订阅。" + ), + } + else: + raise clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY) db.commit() _clear_active_subscription() + clear_ume_subscription_lost_flag() request_ws_reconnect() - append_ws_log(f"manual cancel subscription id={sub_id or '(none)'}", subscription_id=sub_id) - return get_subscription_status() + append_ws_log( + f"manual cancel subscription id={sub_id or '(none)'}", + subscription_id=sub_id, + ) + st = get_subscription_status() + return { + **st, + "ok": True, + "cleared_local": True, + "ume_already_missing": ume_already_missing, + } def _run_ws_session( @@ -357,14 +464,25 @@ def _run_ws_session( db.close() def _on_error(_ws: Any, error: Any) -> None: - detail = f"ws error: {str(error)[:200]}" - _notify_ws_connection( - "error", - detail=detail, - subscription_id=subscription_id, - on_status=on_status, - log_level="error", - ) + err_text = str(error) + detail = f"ws error: {err_text[:200]}" + if is_ume_subscription_missing_error(err_text): + mark_ume_subscription_lost(err_text) + _notify_ws_connection( + "subscription_lost", + detail=detail, + subscription_id=subscription_id, + on_status=on_status, + log_level="warning", + ) + else: + _notify_ws_connection( + "error", + detail=detail, + subscription_id=subscription_id, + on_status=on_status, + log_level="error", + ) def _on_close(_ws: Any, close_status_code: Any, close_msg: Any) -> None: detail = f"ws closed code={close_status_code} msg={str(close_msg or '')[:120]}" @@ -461,6 +579,18 @@ def run_alarm_ws_consumer_loop( time.sleep(1.0) continue + if is_ume_subscription_lost(): + with _ume_lost_lock: + lost_detail = str(_ume_subscription_lost_reason or "") + _loop_status( + "subscription_lost", + detail=lost_detail or "UME subscription missing on server", + log_level="warning", + ) + _ws_wake_event.wait(timeout=5.0) + _ws_wake_event.clear() + continue + subscription_id, wss_uri = get_active_subscription() if not subscription_id or not wss_uri: _loop_status("no_subscription") diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 303abcf..4d070ab 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -18,11 +18,14 @@ from netx_api.ume_alarm_subscription_store import clear_subscription, load_subsc from netx_api.ume_alarm_ws import ( append_ws_log, cancel_alarm_subscription_manual, + clear_local_alarm_subscription_manual, establish_alarm_subscription_manual, get_subscription_status, get_ws_connection_status, get_ws_logs, + is_ume_subscription_missing_error, load_persisted_subscription, + mark_ume_subscription_lost, parse_subscription_id_from_already_exists_error, process_alarm_notification, ) @@ -739,6 +742,53 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertTrue(changed) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-1")) + def test_is_ume_subscription_missing_error(self): + err = ( + "ConnectFailed: Server responded with a non-101 status: 503 " + "status-reason: Subscription not found or overtime! Please establish again" + ) + self.assertTrue(is_ume_subscription_missing_error(err)) + self.assertTrue( + is_ume_subscription_missing_error( + 'ume_request_failed:DELETE ... 400:{"errorInfo":"subscription not exist!"}' + ) + ) + + def test_cancel_returns_needs_cleanup_when_ume_missing(self): + from time import time as _time + + ume_alarm_ws._clear_active_subscription() + ume_alarm_ws.clear_ume_subscription_lost_flag() + 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 + save_subscription(db, subscription_id="sub-gone", wss_uri="wss://ume.local/stream/sub-gone", topic="ALARM") + db.commit() + ume_alarm_ws._set_active_subscription("sub-gone", "wss://ume.local/stream/sub-gone") + + def _fail_delete(_id: str) -> None: + raise RuntimeError('ume_request_failed:DELETE ... 400:{"errorInfo":"subscription not exist!"}') + + client.delete_alarm_subscription = _fail_delete # type: ignore[method-assign] + client.refresh_if_needed = lambda: "token-1" # type: ignore[method-assign] + + st = cancel_alarm_subscription_manual(client, db) + self.assertFalse(st.get("ok", True)) + self.assertTrue(st.get("needs_local_cleanup")) + self.assertTrue(st.get("active")) + self.assertTrue(ume_alarm_ws.is_ume_subscription_lost()) + + st2 = cancel_alarm_subscription_manual(client, db, force_clear_local=True) + self.assertTrue(st2.get("ok")) + self.assertFalse(st2.get("active")) + self.assertFalse(ume_alarm_ws.is_ume_subscription_lost()) + db.close() + def test_ws_connection_status_not_alarm_action(self): from netx_api import ume_alarm_ws as ws_mod diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index f4fcf21..ab22b8b 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -6,6 +6,7 @@ import { fetchUmeCurrentAlarms, fetchUmeNe, cancelUmeAlarmSubscription, + clearLocalUmeAlarmSubscription, establishUmeAlarmSubscription, fetchUmeAlarmSubscriptionStatus, fetchUmeSyncStatus, @@ -89,8 +90,13 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { staleTime: 3000, refetchInterval: 5000, }); + const confirmClearLocalSubscription = (hint?: string) => + window.confirm( + `${hint || "UME 侧告警订阅已不存在或已过期(服务器订阅已丢失)。"}\n\n是否清除本地订阅记录?清除后可点击「建立告警订阅」重新订阅。`, + ); + const subscriptionEstablishMutation = useMutation({ - mutationFn: establishUmeAlarmSubscription, + mutationFn: (opts?: { forceReestablish?: boolean }) => establishUmeAlarmSubscription(opts), onMutate: () => setSubscriptionOpError(""), onSuccess: async (res) => { if (!res?.active) { @@ -114,17 +120,47 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { toastError?.(msg); }, }); + const subscriptionClearLocalMutation = useMutation({ + mutationFn: clearLocalUmeAlarmSubscription, + 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 subscriptionCancelMutation = useMutation({ - mutationFn: cancelUmeAlarmSubscription, + mutationFn: (opts?: { forceClearLocal?: boolean }) => cancelUmeAlarmSubscription(opts), onMutate: () => setSubscriptionOpError(""), onSuccess: async (res) => { + if (res?.needs_local_cleanup) { + const msg = + res.message || + "UME 侧告警订阅已丢失。请确认是否清除本地订阅记录后重新订阅。"; + if (confirmClearLocalSubscription(msg)) { + subscriptionCancelMutation.mutate({ forceClearLocal: true }); + return; + } + setSubscriptionOpError(msg); + toastError?.(msg); + await queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] }); + await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] }); + return; + } if (res?.active) { const msg = "取消订阅失败:UME 侧订阅可能仍存在,请查看错误详情后重试"; setSubscriptionOpError(msg); toastError?.(msg); } else { setSubscriptionOpError(""); - toastOk?.("告警订阅已取消"); + toastOk?.(res?.ume_already_missing ? "本地已清除(UME 侧订阅此前已不存在)" : "告警订阅已取消"); } await queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] }); await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] }); @@ -204,21 +240,33 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { wsConn?.label || (wsPaused ? "已暂停" : wsConsumer?.last_error || wsConsumer?.status || "-"); const wsPillLevel = - wsPaused || wsState === "paused" + serverSubLost || wsState === "subscription_lost" ? "warn" - : wsState === "connected" - ? "up" - : wsState === "connecting" || wsState === "reconnecting" || wsState === "waiting_token" - ? "warn" - : wsState === "no_subscription" || wsState === "disconnected" || wsState === "error" || wsState === "init" - ? "down" - : String(wsConsumer?.last_error || "").includes("connected") - ? "up" - : "unknown"; + : wsPaused || wsState === "paused" + ? "warn" + : wsState === "connected" + ? "up" + : wsState === "connecting" || wsState === "reconnecting" || wsState === "waiting_token" + ? "warn" + : wsState === "no_subscription" || wsState === "disconnected" || wsState === "error" || wsState === "init" + ? "down" + : String(wsConsumer?.last_error || "").includes("connected") + ? "up" + : "unknown"; const wsLogs = [...(subscriptionStatusQuery.data?.ws_logs ?? alarmSub.ws_logs ?? [])].reverse(); const subscriptionActive = Boolean(alarmSub.active); + const serverSubLost = Boolean( + subscriptionStatusQuery.data?.server_subscription_lost ?? alarmSub.server_subscription_lost, + ); + const serverSubLostReason = String( + subscriptionStatusQuery.data?.server_subscription_lost_reason ?? + alarmSub.server_subscription_lost_reason ?? + "", + ); const subPending = - subscriptionEstablishMutation.isPending || subscriptionCancelMutation.isPending; + subscriptionEstablishMutation.isPending || + subscriptionCancelMutation.isPending || + subscriptionClearLocalMutation.isPending; const runtimeTaskMutation = useMutation({ mutationFn: async (vars: { task: string; action: "pause" | "resume" }) => @@ -336,22 +384,72 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { ) : null} + {serverSubLost ? ( +