diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py index 224a911..0004f46 100644 --- a/netx_api/ume_alarm_ws.py +++ b/netx_api/ume_alarm_ws.py @@ -2,6 +2,7 @@ from __future__ import annotations import json import logging +import re import ssl import threading import time @@ -22,6 +23,14 @@ from .ume_sync_service import apply_alarm_to_current, extract_alarm_from_notific _ws_log = logging.getLogger("netx.ume.alarm_ws") +_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})") + + +def parse_subscription_id_from_already_exists_error(message: str) -> str: + """Extract subscription id from UME 400 'topic subscription already exist, id:...'.""" + m = _ORPHAN_SUB_ID_RE.search(str(message or "")) + return m.group(1).strip() if m else "" + _shutdown_lock = threading.Lock() _subscription_lock = threading.Lock() _ws_wake_event = threading.Event() @@ -146,7 +155,15 @@ def establish_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[ 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) + try: + sub_id, uri = client.establish_alarm_subscription(topic=topic) + except RuntimeError as exc: + orphan_id = parse_subscription_id_from_already_exists_error(str(exc)) + if not orphan_id: + raise + _ws_log.warning("establish: UME reports existing subscription id=%s, deleting then retry", orphan_id) + _delete_subscription_on_ume(client, orphan_id) + sub_id, uri = client.establish_alarm_subscription(topic=topic) save_subscription(db, subscription_id=sub_id, wss_uri=uri, topic=topic) db.commit() @@ -166,12 +183,8 @@ def cancel_alarm_subscription_manual(client: UMEClient, db: Session) -> dict[str 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) + client.refresh_if_needed() + _delete_subscription_on_ume(client, sub_id) clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY) db.commit() diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index 3825f50..c6afdc5 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -622,17 +622,12 @@ class UMEClient: return sub_id, uri def delete_alarm_subscription(self, subscription_id: str) -> None: + """DELETE delete-subscription with body {\"input\": {\"id\": ...}} (same headers as alarm REST).""" 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 + self.request_json("DELETE", self.notification_delete_path, body=body) def has_valid_token(self) -> bool: """True when a non-expired token is present in memory (call _sync_token_from_store first).""" diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 3c689b2..e81d194 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -19,6 +19,7 @@ from netx_api.ume_alarm_ws import ( establish_alarm_subscription_manual, get_subscription_status, load_persisted_subscription, + parse_subscription_id_from_already_exists_error, process_alarm_notification, ) from netx_api.ume_sync_service import ( @@ -187,9 +188,9 @@ class UMEClientTests(unittest.TestCase): seen["body"] = body return ({}, None) - client._request_json_with_current_token = _fake_request # type: ignore[method-assign] + client.request_json = _fake_request # type: ignore[method-assign] client.delete_alarm_subscription("sub-to-delete") - self.assertEqual(seen["method"], "POST") + self.assertEqual(seen["method"], "DELETE") self.assertEqual(seen["body"], {"input": {"id": "sub-to-delete"}}) @@ -701,6 +702,82 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertTrue(changed) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-1")) + def test_parse_subscription_id_from_already_exists_error(self): + msg = ( + 'ume_request_failed:400:{"error":{"errorInfo":"topic subscription already exist, ' + 'id:087f3544-0171-49c9-9aea-2ac535b3deaa, uri:wss://10.0.0.1:18014/restconf/stream/087f3544"}}' + ) + self.assertEqual( + parse_subscription_id_from_already_exists_error(msg), + "087f3544-0171-49c9-9aea-2ac535b3deaa", + ) + + def test_cancel_subscription_fails_without_clearing_db(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 + save_subscription(db, subscription_id="sub-keep", wss_uri="wss://ume.local/stream/sub-keep", topic="ALARM") + db.commit() + ume_alarm_ws._set_active_subscription("sub-keep", "wss://ume.local/stream/sub-keep") + + def _fail_delete(_id: str) -> None: + raise RuntimeError("ume_request_failed:500:delete failed") + + client.delete_alarm_subscription = _fail_delete # type: ignore[method-assign] + client.refresh_if_needed = lambda: "token-1" # type: ignore[method-assign] + + with self.assertRaises(RuntimeError): + cancel_alarm_subscription_manual(client, db) + loaded = load_subscription(db) + self.assertIsNotNone(loaded) + self.assertEqual(loaded[0], "sub-keep") + self.assertTrue(get_subscription_status()["active"]) + db.close() + + def test_establish_recovers_when_ume_reports_already_exists(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 + deleted: list[str] = [] + calls = {"n": 0} + + def _fake_establish(*, topic=None): + calls["n"] += 1 + if calls["n"] == 1: + raise RuntimeError( + "ume_request_failed:400:{\"error\":{\"errorInfo\":\"topic subscription already exist, " + "id:087f3544-0171-49c9-9aea-2ac535b3deaa\"}}" + ) + return ("sub-new", "wss://ume.local/stream/sub-new") + + def _fake_delete(sub_id: str) -> None: + deleted.append(sub_id) + + client.establish_alarm_subscription = _fake_establish # type: ignore[method-assign] + client.delete_alarm_subscription = _fake_delete # type: ignore[method-assign] + + st = establish_alarm_subscription_manual(client, db) + self.assertEqual(deleted, ["087f3544-0171-49c9-9aea-2ac535b3deaa"]) + self.assertEqual(st["subscription_id"], "sub-new") + self.assertTrue(st["active"]) + db.close() + def test_subscription_store_and_manual_establish(self): from time import time as _time diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index 52f71b7..b1f2b6d 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -117,9 +117,15 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { const subscriptionCancelMutation = useMutation({ mutationFn: cancelUmeAlarmSubscription, onMutate: () => setSubscriptionOpError(""), - onSuccess: async () => { - setSubscriptionOpError(""); - toastOk?.("告警订阅已取消"); + onSuccess: async (res) => { + if (res?.active) { + const msg = "取消订阅失败:UME 侧订阅可能仍存在,请查看错误详情后重试"; + setSubscriptionOpError(msg); + toastError?.(msg); + } else { + setSubscriptionOpError(""); + toastOk?.("告警订阅已取消"); + } await queryClient.invalidateQueries({ queryKey: ["umeAlarmSubscription"] }); await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] }); },