diff --git a/.env.example b/.env.example index 8faa188..116a311 100644 --- a/.env.example +++ b/.env.example @@ -21,6 +21,9 @@ 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=18000 +NETX_UME_SYNC_ALARMS_CURRENT_SKIP_WHEN_WS=true +NETX_UME_STARTUP_SYNC_ALARMS_BEFORE_WS=true +NETX_UME_ALARM_CLEARED_TOMBSTONE_S=300 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 diff --git a/netx_api/config.py b/netx_api/config.py index d042d7e..010f94a 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -42,6 +42,9 @@ class Settings(BaseSettings): ume_keepalive_renew_before_s: int = 900 ume_sync_alarms_current_enabled: bool = True ume_sync_alarms_current_interval_s: int = 18000 + ume_sync_alarms_current_skip_when_ws: bool = True + ume_startup_sync_alarms_before_ws: bool = True + ume_alarm_cleared_tombstone_s: int = 300 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" diff --git a/netx_api/main.py b/netx_api/main.py index 0d0567a..42f434f 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -35,12 +35,16 @@ from .models import ImportJob from .parser_config import load_parser_config from .ume_client import UMEClient from .ume_alarm_ws import ( + begin_startup_alarm_sync_gate, cancel_alarm_subscription_manual, clear_local_alarm_subscription_manual, + complete_startup_alarm_sync_gate, establish_alarm_subscription_manual, + get_alarms_coordination_status, get_subscription_status, get_ws_connection_status, get_ws_logs, + is_wss_active_for_current_alarms, load_persisted_subscription, request_ws_reconnect, shutdown_ws_consumer, @@ -274,6 +278,51 @@ def _fail_stale_running_sync_jobs_on_startup() -> None: db.close() +def _run_startup_alarm_sync_before_ws() -> None: + """REST-sync current alarms once on boot before WSS connects (avoids stale/reconcile races).""" + ume_url = str(getattr(settings, "ume_base_url", "") or "").strip() + alarms_enabled = bool(getattr(settings, "ume_sync_alarms_current_enabled", True)) + ws_enabled = bool(getattr(settings, "ume_alarm_ws_enabled", True)) + before_ws = bool(getattr(settings, "ume_startup_sync_alarms_before_ws", True)) + + if not (before_ws and ws_enabled and alarms_enabled and ume_url): + complete_startup_alarm_sync_gate() + return + + begin_startup_alarm_sync_gate() + try: + _schedule_log.info("startup: syncing current alarms before WSS") + _set_runtime_task( + "alarms_current_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error="启动:正在同步当前告警(完成后连接 WSS)…", + ) + db = SessionLocal() + try: + client = _ume_client() + sync_alarms_current(db, client, trigger_mode="schedule", wss_active=False) + _schedule_log.info("startup: current alarms sync completed, WSS may connect") + _set_runtime_task( + "alarms_current_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error="", + ) + finally: + db.close() + except Exception as exc: + _schedule_log.exception("startup: current alarms sync before WSS failed: %s", exc) + _set_runtime_task( + "alarms_current_auto_sync", + status="error", + last_run_at=datetime.now(timezone.utc), + last_error=str(exc)[:240], + ) + finally: + complete_startup_alarm_sync_gate() + + def _sleep_or_until_paused(task_id: str, total_s: float) -> None: """Sleep up to total_s wall seconds; honor pause; wake early on resume (debounce interrupt).""" deadline = time.time() + max(0.0, float(total_s)) @@ -717,6 +766,11 @@ def on_startup() -> None: last_run_at=datetime.now(timezone.utc), last_error=f"startup_thread_init_failed: {str(exc)[:180]}", ) + try: + _run_startup_alarm_sync_before_ws() + except Exception as exc: + _schedule_log.exception("startup: alarm sync before WSS failed: %s", exc) + complete_startup_alarm_sync_gate() try: if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)): alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000) @@ -732,6 +786,18 @@ def on_startup() -> None: if _runtime_is_paused("alarms_current_auto_sync"): time.sleep(1) continue + if ( + bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True)) + and is_wss_active_for_current_alarms() + ): + _set_runtime_task( + "alarms_current_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error="WSS 实时接收中,已跳过 REST 同步", + ) + time.sleep(max(30, min(alarms_interval_s, 300))) + continue _maybe_wait_for_sync_interval( task_id="alarms_current_auto_sync", domain="alarms_current", @@ -938,6 +1004,7 @@ def ume_alarm_subscription_status(limit: int = 80) -> dict[str, Any]: return { "ok": True, **st, + **get_alarms_coordination_status(), "ws_connection": get_ws_connection_status(), "ws_consumer_status": str(ws_task.get("status") or ""), "ws_consumer_last_error": str(ws_task.get("last_error") or ""), @@ -1017,7 +1084,22 @@ def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db } ) if "alarms" in domain_set or "alarms_current" in domain_set: - job, batch = sync_alarms_current(db, client, trigger_mode=trigger_mode) + paused_ws_for_sync = False + if is_wss_active_for_current_alarms() and trigger_mode == "manual": + _runtime_pause_task("alarms_current_ws_consumer") + request_ws_reconnect() + paused_ws_for_sync = True + try: + job, batch = sync_alarms_current( + db, + client, + trigger_mode=trigger_mode, + wss_active=is_wss_active_for_current_alarms(), + ) + finally: + if paused_ws_for_sync: + _runtime_resume_task("alarms_current_ws_consumer") + request_ws_reconnect() out["jobs"].append( { "domain": "alarms_current", diff --git a/netx_api/ume_alarm_ws.py b/netx_api/ume_alarm_ws.py index 1877441..0a8c9a1 100644 --- a/netx_api/ume_alarm_ws.py +++ b/netx_api/ume_alarm_ws.py @@ -42,6 +42,22 @@ _WS_LOG_LOCK = threading.Lock() _WS_LOG_ENTRIES: deque[dict[str, Any]] = deque(maxlen=200) _WS_LOG_MAX_RETURN = 100 +# Blocks WSS connect until startup REST sync of current alarms completes (see main.on_startup). +_STARTUP_ALARM_SYNC_GATE = threading.Event() +_STARTUP_ALARM_SYNC_GATE.set() + + +def begin_startup_alarm_sync_gate() -> None: + _STARTUP_ALARM_SYNC_GATE.clear() + + +def complete_startup_alarm_sync_gate() -> None: + _STARTUP_ALARM_SYNC_GATE.set() + + +def is_startup_alarm_sync_pending() -> bool: + return not _STARTUP_ALARM_SYNC_GATE.is_set() + _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, ...] = ( @@ -283,6 +299,29 @@ def get_subscription_status() -> dict[str, Any]: } +def is_wss_active_for_current_alarms() -> bool: + """True when WSS subscription is active and should own ume_alarms_current updates.""" + if not bool(getattr(settings, "ume_alarm_ws_enabled", True)): + return False + if is_ume_subscription_lost(): + return False + return bool(get_subscription_status().get("active")) + + +def get_current_alarms_mode() -> str: + return "wss" if is_wss_active_for_current_alarms() else "rest" + + +def get_alarms_coordination_status() -> dict[str, Any]: + wss_active = is_wss_active_for_current_alarms() + skip_when_ws = bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True)) + return { + "current_alarms_mode": get_current_alarms_mode(), + "wss_active_for_current_alarms": wss_active, + "scheduled_sync_skipped": bool(wss_active and skip_when_ws), + } + + 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) @@ -579,6 +618,11 @@ def run_alarm_ws_consumer_loop( time.sleep(1.0) continue + if not _STARTUP_ALARM_SYNC_GATE.is_set(): + _loop_status("waiting_startup_alarm_sync", log_level="info") + _STARTUP_ALARM_SYNC_GATE.wait(timeout=1.0) + continue + if is_ume_subscription_lost(): with _ume_lost_lock: lost_detail = str(_ume_subscription_lost_reason or "") diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index b469b25..10fd9a9 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -4,12 +4,17 @@ import json import logging import hashlib import re +import threading +import time from datetime import datetime, timezone from typing import Any _sync_log = logging.getLogger("netx.ume.sync") +from sqlalchemy import func from sqlalchemy import text as sql_text +from sqlalchemy.dialects.postgresql import insert as pg_insert +from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session from .config import settings @@ -196,6 +201,127 @@ def _is_alarm_cleared(alarm: dict[str, Any]) -> bool: return text in {"true", "1", "yes"} +_cleared_tombstone_lock = threading.Lock() +_cleared_tombstones: dict[str, float] = {} + + +def _mark_alarm_cleared_tombstone(alarm_key: str) -> None: + key = str(alarm_key or "").strip() + if not key: + return + ttl_s = max(60, int(getattr(settings, "ume_alarm_cleared_tombstone_s", 300) or 300)) + expires = time.time() + ttl_s + with _cleared_tombstone_lock: + _cleared_tombstones[key] = expires + if len(_cleared_tombstones) > 50000: + now = time.time() + stale = [k for k, exp in _cleared_tombstones.items() if exp <= now] + for k in stale[:10000]: + _cleared_tombstones.pop(k, None) + + +def _is_alarm_cleared_tombstone(alarm_key: str) -> bool: + key = str(alarm_key or "").strip() + if not key: + return False + now = time.time() + with _cleared_tombstone_lock: + exp = _cleared_tombstones.get(key) + if exp is None: + return False + if exp <= now: + _cleared_tombstones.pop(key, None) + return False + return True + + +def _alarm_row_from_norm(key: str, norm: dict[str, Any], *, touch_ts: datetime, first_seen_at: datetime) -> dict[str, Any]: + return { + "alarm_key": key, + "ne_id": _s(_derive_ne_id_from_alarm(norm)), + "host_name": "", + "object_name": _s(_pick(norm, "objectName", "object-name")), + "event_type": _s(_pick(norm, "eventType", "event-type")), + "native_probable_cause": _s(_pick(norm, "nativeProbableCause", "native-probable-cause")), + "perceived_severity": _s(_pick(norm, "perceivedSeverity", "perceived-severity")), + "is_cleared": _s(_pick(norm, "isCleared", "is-cleared")), + "time_created": _s(_pick(norm, "timeCreated", "time-created")), + "root_cause_alarm_indication": _s( + _pick(norm, "rootCauseAlarmIndication", "root-cause-alarm-indication") + ), + "first_seen_at": first_seen_at, + "last_seen_at": touch_ts, + "raw_json": json.dumps(norm, ensure_ascii=False, default=str), + } + + +def _apply_row_to_model(db: Session, existing: UmeAlarmCurrent, norm: dict[str, Any], *, touch_ts: datetime) -> None: + existing.ne_id = _s(_derive_ne_id_from_alarm(norm)) + 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") + ) + prev_seen = existing.last_seen_at + if prev_seen is None or touch_ts >= prev_seen: + existing.last_seen_at = touch_ts + existing.raw_json = json.dumps(norm, ensure_ascii=False, default=str) + + +def _upsert_alarm_current(db: Session, key: str, norm: dict[str, Any], *, touch_ts: datetime) -> tuple[str, bool]: + bind = db.get_bind() + dialect = str(getattr(getattr(bind, "dialect", None), "name", "") or "").lower() + existing = db.get(UmeAlarmCurrent, key) + if existing is not None: + _apply_row_to_model(db, existing, norm, touch_ts=touch_ts) + existing.host_name = _lookup_host_name(db, existing.ne_id) + return "updated", True + + if dialect == "postgresql": + row = _alarm_row_from_norm(key, norm, touch_ts=touch_ts, first_seen_at=touch_ts) + row["host_name"] = _lookup_host_name(db, row["ne_id"]) + ins = pg_insert(UmeAlarmCurrent).values(**row) + excluded = ins.excluded + stmt = ins.on_conflict_do_update( + index_elements=[UmeAlarmCurrent.alarm_key], + set_={ + "ne_id": excluded.ne_id, + "host_name": excluded.host_name, + "object_name": excluded.object_name, + "event_type": excluded.event_type, + "native_probable_cause": excluded.native_probable_cause, + "perceived_severity": excluded.perceived_severity, + "is_cleared": excluded.is_cleared, + "time_created": excluded.time_created, + "root_cause_alarm_indication": excluded.root_cause_alarm_indication, + "last_seen_at": func.greatest(UmeAlarmCurrent.last_seen_at, excluded.last_seen_at), + "raw_json": excluded.raw_json, + }, + ) + db.execute(stmt) + return "inserted", True + + try: + with db.begin_nested(): + model = UmeAlarmCurrent(alarm_key=key, first_seen_at=touch_ts) + db.add(model) + db.flush() + _apply_row_to_model(db, model, norm, touch_ts=touch_ts) + model.host_name = _lookup_host_name(db, model.ne_id) + return "inserted", True + except IntegrityError: + existing = db.get(UmeAlarmCurrent, key) + if existing is None: + return "skipped", False + _apply_row_to_model(db, existing, norm, touch_ts=touch_ts) + existing.host_name = _lookup_host_name(db, existing.ne_id) + return "updated", True + + 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): @@ -223,7 +349,13 @@ def extract_alarm_from_notification(payload: dict[str, Any]) -> dict[str, Any] | 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]: +def apply_alarm_to_current( + db: Session, + alarm: dict[str, Any], + *, + touch_ts: datetime, + source: str = "", +) -> tuple[str, bool]: """ Apply one alarm to ume_alarms_current. Returns (action, changed) where action is inserted|updated|deleted|skipped. @@ -238,32 +370,19 @@ def apply_alarm_to_current(db: Session, alarm: dict[str, Any], *, touch_ts: date return "skipped", False existing = db.get(UmeAlarmCurrent, key) if existing is None: + _mark_alarm_cleared_tombstone(key) return "deleted", False db.delete(existing) + _mark_alarm_cleared_tombstone(key) 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 + if not key: + return "skipped", False + if str(source or "").strip().lower() == "rest" and _is_alarm_cleared_tombstone(key): + return "skipped", False + + return _upsert_alarm_current(db, key, norm, touch_ts=touch_ts) def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: @@ -471,12 +590,34 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " return job +def _reconcile_stale_current_alarms( + db: Session, + *, + sync_batch_ts: datetime, + seen_keys: set[str], + wss_active: bool, +) -> int: + """Remove local current alarms missing from REST snapshot. + + When WSS is active, skip deletes (WSS may have keys not yet in REST); manual sync only upserts. + """ + del seen_keys + if wss_active: + return 0 + return int( + db.query(UmeAlarmCurrent) + .filter(UmeAlarmCurrent.last_seen_at < sync_batch_ts) + .delete(synchronize_session=False) + ) + + def _sync_alarms_common( db: Session, client: UMEClient, *, is_uncleared: bool, trigger_mode: str, + wss_active: bool = False, ) -> tuple[UmeSyncJob, UmeAlarmBatch]: domain = "alarms_history" if is_uncleared else "alarms_current" job = _build_sync_job(domain, trigger_mode) @@ -498,6 +639,8 @@ def _sync_alarms_common( pulled = inserted = updated = 0 deleted_stale_current = 0 host_names_backfilled = 0 + reconcile_mode = "" + seen_keys: set[str] = set() paging_mode = "marker" paging_note = "" page_no = 0 @@ -506,6 +649,7 @@ def _sync_alarms_common( graceful_end_by_iterator_error = False warnings: list[str] = [] meta: dict[str, Any] = {} + sync_batch_ts = _utc_now_naive() try: limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = max(1, limit_max) @@ -545,13 +689,22 @@ def _sync_alarms_common( iterator_500_as_end=iterator_500_as_end, ) sync_batch_ts = _utc_now_naive() + seen_keys = set() for rows in pages: pulled += len(rows) for alarm in rows: 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) + key = _alarm_key(alarm) + if key: + seen_keys.add(key) + action, _changed = apply_alarm_to_current( + db, + alarm, + touch_ts=sync_batch_ts, + source="rest", + ) if action == "inserted": inserted += 1 elif action == "updated": @@ -564,11 +717,15 @@ def _sync_alarms_common( paging_note = str(meta.get("paging_note") or "") warnings = [str(x) for x in (meta.get("warnings") or []) if str(x)] + reconcile_mode = "full" if not is_uncleared and _snapshot_reconcile_ok(meta): - deleted_stale_current = int( - db.query(UmeAlarmCurrent) - .filter(UmeAlarmCurrent.last_seen_at < sync_batch_ts) - .delete(synchronize_session=False) + if wss_active: + reconcile_mode = "upsert_only" + deleted_stale_current = _reconcile_stale_current_alarms( + db, + sync_batch_ts=sync_batch_ts, + seen_keys=seen_keys, + wss_active=wss_active, ) alarm_model = UmeAlarmHistory if is_uncleared else UmeAlarmCurrent @@ -593,6 +750,9 @@ def _sync_alarms_common( "deleted_stale_current_alarms": int(deleted_stale_current), "host_names_backfilled": int(host_names_backfilled), "current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta), + "reconcile_mode": reconcile_mode if not is_uncleared else "", + "wss_active_during_sync": bool(wss_active) if not is_uncleared else False, + "seen_keys_count": len(seen_keys) if not is_uncleared else 0, }, ensure_ascii=False, ) @@ -600,6 +760,7 @@ def _sync_alarms_common( job.status = "done" except Exception as exc: msg = str(exc)[:1024] + reconcile_mode = "failed" batch.status = "failed" batch.error_message = msg batch.ended_at = _utc_now_naive() @@ -625,6 +786,9 @@ def _sync_alarms_common( "deleted_stale_current_alarms": int(deleted_stale_current), "host_names_backfilled": int(host_names_backfilled), "current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta), + "reconcile_mode": reconcile_mode if not is_uncleared else "", + "wss_active_during_sync": bool(wss_active) if not is_uncleared else False, + "seen_keys_count": len(seen_keys) if not is_uncleared else 0, }, ensure_ascii=False, ) @@ -634,8 +798,24 @@ def _sync_alarms_common( return job, batch -def sync_alarms_current(db: Session, client: UMEClient, *, trigger_mode: str = "manual") -> tuple[UmeSyncJob, UmeAlarmBatch]: - return _sync_alarms_common(db, client, is_uncleared=False, trigger_mode=trigger_mode) +def sync_alarms_current( + db: Session, + client: UMEClient, + *, + trigger_mode: str = "manual", + wss_active: bool | None = None, +) -> tuple[UmeSyncJob, UmeAlarmBatch]: + if wss_active is None: + from .ume_alarm_ws import is_wss_active_for_current_alarms + + wss_active = is_wss_active_for_current_alarms() + return _sync_alarms_common( + db, + client, + is_uncleared=False, + trigger_mode=trigger_mode, + wss_active=bool(wss_active), + ) def sync_alarms_history_full( diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 4d070ab..20f3734 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -388,7 +388,7 @@ class UmeSyncServiceTests(unittest.TestCase): ) self.db.commit() - sync_alarms_current(self.db, _COne(), trigger_mode="manual") + sync_alarms_current(self.db, _COne(), trigger_mode="manual", wss_active=False) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-KEEP")) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-GONE")) @@ -421,7 +421,7 @@ class UmeSyncServiceTests(unittest.TestCase): ) self.db.commit() - sync_alarms_current(self.db, _CPartial(), trigger_mode="manual") + sync_alarms_current(self.db, _CPartial(), trigger_mode="manual", wss_active=False) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-NEW")) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-STALE")) @@ -458,11 +458,11 @@ class UmeSyncServiceTests(unittest.TestCase): ) self.db.commit() - job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual", wss_active=False) self.assertEqual(job1.status, "done") self.assertEqual(job1.inserted_count, 1) - job2, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + job2, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual", wss_active=False) self.assertEqual(job2.status, "done") self.assertEqual(job2.updated_count, 1) @@ -502,7 +502,7 @@ class UmeSyncServiceTests(unittest.TestCase): svc.settings.ume_page_size = 2 svc.settings.ume_max_pages = 10 try: - job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual", wss_active=False) finally: svc.settings.ume_page_size = old_page_size svc.settings.ume_max_pages = old_max_pages @@ -538,7 +538,7 @@ class UmeSyncServiceTests(unittest.TestCase): svc.settings.ume_marker_max_pages = 10 svc.settings.ume_iterator_500_as_end = True try: - job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual", wss_active=False) finally: svc.settings.ume_marker_page_limit = old_page_size svc.settings.ume_marker_max_pages = old_max_pages @@ -568,7 +568,7 @@ class UmeSyncServiceTests(unittest.TestCase): ) c = _C() - job, _ = sync_alarms_current(self.db, c, trigger_mode="manual") + job, _ = sync_alarms_current(self.db, c, trigger_mode="manual", wss_active=False) self.assertEqual(job.status, "done") self.assertEqual(job.pulled_count, 1) self.assertEqual(c.calls, 1) @@ -594,7 +594,7 @@ class UmeSyncServiceTests(unittest.TestCase): ) c = _C() - job, _ = sync_alarms_current(self.db, c, trigger_mode="manual") + job, _ = sync_alarms_current(self.db, c, trigger_mode="manual", wss_active=False) self.assertEqual(job.status, "done") self.assertEqual(job.pulled_count, 1) self.assertEqual(c.calls, 1) @@ -934,6 +934,119 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertIsNone(load_subscription(db)) db.close() + def test_sync_current_alarms_wss_active_skips_stale_delete(self): + from datetime import datetime, timedelta + + class _COne: + def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = True + + return ( + [ + { + "alarmKey": "AK-REST", + "ne-id": "NE-1", + "perceivedSeverity": "major", + "isCleared": "false", + }, + ], + _D(), + ) + + ws_touch = datetime.utcnow() + self.db.add( + UmeAlarmCurrent( + alarm_key="AK-WS-ONLY", + ne_id="NE-2", + perceived_severity="minor", + is_cleared="false", + last_seen_at=ws_touch, + ) + ) + self.db.add( + UmeAlarmCurrent( + alarm_key="AK-STALE", + ne_id="NE-X", + perceived_severity="minor", + is_cleared="false", + last_seen_at=ws_touch - timedelta(hours=1), + ) + ) + self.db.commit() + + job, batch = sync_alarms_current(self.db, _COne(), trigger_mode="manual", wss_active=True) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-REST")) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-WS-ONLY")) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-STALE")) + details = json.loads(batch.raw_json or "{}") + self.assertEqual(int(details.get("deleted_stale_current_alarms") or 0), 0) + self.assertEqual(str(details.get("reconcile_mode") or ""), "upsert_only") + self.assertEqual(job.status, "done") + + def test_apply_alarm_rest_skips_tombstone_after_wss_clear(self): + from datetime import datetime + + touch = datetime.utcnow() + self.db.add( + UmeAlarmCurrent( + alarm_key="AK-CLEARED-2", + ne_id="NE-1", + is_cleared="false", + ) + ) + self.db.commit() + action, changed = apply_alarm_to_current( + self.db, + {"alarmKey": "AK-CLEARED-2", "isCleared": True}, + touch_ts=touch, + ) + self.db.commit() + self.assertEqual(action, "deleted") + self.assertTrue(changed) + action2, changed2 = apply_alarm_to_current( + self.db, + {"alarmKey": "AK-CLEARED-2", "isCleared": False, "perceivedSeverity": "major"}, + touch_ts=touch, + source="rest", + ) + self.assertEqual(action2, "skipped") + self.assertFalse(changed2) + self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-CLEARED-2")) + + def test_apply_alarm_concurrent_upsert_no_integrity_error(self): + from datetime import datetime + + touch = datetime.utcnow() + alarm = { + "alarmKey": "AK-DUP", + "ne-id": "NE-1", + "perceivedSeverity": "major", + "isCleared": "false", + } + a1, _ = apply_alarm_to_current(self.db, alarm, touch_ts=touch) + self.db.commit() + a2, _ = apply_alarm_to_current(self.db, alarm, touch_ts=touch, source="rest") + self.db.commit() + self.assertEqual(a1, "inserted") + self.assertIn(a2, ("updated", "inserted")) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-DUP")) + + def test_is_wss_active_for_current_alarms(self): + from netx_api.ume_alarm_ws import ( + _clear_active_subscription, + _set_active_subscription, + clear_ume_subscription_lost_flag, + is_wss_active_for_current_alarms, + ) + + _clear_active_subscription() + clear_ume_subscription_lost_flag() + self.assertFalse(is_wss_active_for_current_alarms()) + _set_active_subscription("sub-x", "wss://ume.local/stream/sub-x") + self.assertTrue(is_wss_active_for_current_alarms()) + def test_process_alarm_notification_via_ws_helper(self): from datetime import datetime diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index de7adbe..67d2c43 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -248,6 +248,11 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { alarmSub.server_subscription_lost_reason ?? "", ); + const currentAlarmsMode = + subscriptionStatusQuery.data?.current_alarms_mode ?? alarmSub.current_alarms_mode ?? "rest"; + const scheduledSyncSkipped = Boolean( + subscriptionStatusQuery.data?.scheduled_sync_skipped ?? alarmSub.scheduled_sync_skipped, + ); const wsPillLevel = serverSubLost || wsState === "subscription_lost" ? "warn" @@ -384,6 +389,11 @@ export function UmePage({ toastOk, toastError }: UmePageProps) { ) : null} + {scheduledSyncSkipped && !serverSubLost ? ( +
+ 当前告警由 WSS 实时维护(模式: {currentAlarmsMode}),定时 REST 全量同步已暂停。需要全量对账时请使用下方「同步当前告警」。 +
+ ) : null} {serverSubLost ? (
UME 服务器侧告警订阅已丢失或已过期,本地记录可能仍显示「已建立」。请确认清除本地记录后重新订阅。 diff --git a/web/src/types.ts b/web/src/types.ts index f59ec92..668da94 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -112,6 +112,9 @@ export type UmeAlarmSubscriptionStatus = { topic?: string; server_subscription_lost?: boolean; server_subscription_lost_reason?: string; + current_alarms_mode?: "wss" | "rest"; + wss_active_for_current_alarms?: boolean; + scheduled_sync_skipped?: boolean; needs_local_cleanup?: boolean; ume_already_missing?: boolean; cleared_local?: boolean;