feat(ume): coordinate WSS with REST sync and sync before WSS on startup

WSS-primary current alarms with upsert/tombstone, skip scheduled REST when WSS active,
safe manual reconcile, startup REST baseline before WebSocket connect.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-27 09:57:33 +08:00
parent aaee32a7a8
commit 6701aee3c7
8 changed files with 476 additions and 38 deletions

View file

@ -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_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list
NETX_UME_SYNC_ALARMS_CURRENT_ENABLED=true NETX_UME_SYNC_ALARMS_CURRENT_ENABLED=true
NETX_UME_SYNC_ALARMS_CURRENT_INTERVAL_S=18000 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_ALARM_WS_ENABLED=true
NETX_UME_NOTIFICATION_ESTABLISH_PATH=/restconf/operations/zte-notifications:establish-subscription 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_DELETE_PATH=/restconf/operations/zte-notifications:delete-subscription

View file

@ -42,6 +42,9 @@ class Settings(BaseSettings):
ume_keepalive_renew_before_s: int = 900 ume_keepalive_renew_before_s: int = 900
ume_sync_alarms_current_enabled: bool = True ume_sync_alarms_current_enabled: bool = True
ume_sync_alarms_current_interval_s: int = 18000 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_alarm_ws_enabled: bool = True
ume_notification_establish_path: str = "/restconf/operations/zte-notifications:establish-subscription" ume_notification_establish_path: str = "/restconf/operations/zte-notifications:establish-subscription"
ume_notification_delete_path: str = "/restconf/operations/zte-notifications:delete-subscription" ume_notification_delete_path: str = "/restconf/operations/zte-notifications:delete-subscription"

View file

@ -35,12 +35,16 @@ from .models import ImportJob
from .parser_config import load_parser_config from .parser_config import load_parser_config
from .ume_client import UMEClient from .ume_client import UMEClient
from .ume_alarm_ws import ( from .ume_alarm_ws import (
begin_startup_alarm_sync_gate,
cancel_alarm_subscription_manual, cancel_alarm_subscription_manual,
clear_local_alarm_subscription_manual, clear_local_alarm_subscription_manual,
complete_startup_alarm_sync_gate,
establish_alarm_subscription_manual, establish_alarm_subscription_manual,
get_alarms_coordination_status,
get_subscription_status, get_subscription_status,
get_ws_connection_status, get_ws_connection_status,
get_ws_logs, get_ws_logs,
is_wss_active_for_current_alarms,
load_persisted_subscription, load_persisted_subscription,
request_ws_reconnect, request_ws_reconnect,
shutdown_ws_consumer, shutdown_ws_consumer,
@ -274,6 +278,51 @@ def _fail_stale_running_sync_jobs_on_startup() -> None:
db.close() 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: 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).""" """Sleep up to total_s wall seconds; honor pause; wake early on resume (debounce interrupt)."""
deadline = time.time() + max(0.0, float(total_s)) 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_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}", last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
) )
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: try:
if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)): if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000) 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"): if _runtime_is_paused("alarms_current_auto_sync"):
time.sleep(1) time.sleep(1)
continue 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( _maybe_wait_for_sync_interval(
task_id="alarms_current_auto_sync", task_id="alarms_current_auto_sync",
domain="alarms_current", domain="alarms_current",
@ -938,6 +1004,7 @@ def ume_alarm_subscription_status(limit: int = 80) -> dict[str, Any]:
return { return {
"ok": True, "ok": True,
**st, **st,
**get_alarms_coordination_status(),
"ws_connection": get_ws_connection_status(), "ws_connection": get_ws_connection_status(),
"ws_consumer_status": str(ws_task.get("status") or ""), "ws_consumer_status": str(ws_task.get("status") or ""),
"ws_consumer_last_error": str(ws_task.get("last_error") 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: 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( out["jobs"].append(
{ {
"domain": "alarms_current", "domain": "alarms_current",

View file

@ -42,6 +42,22 @@ _WS_LOG_LOCK = threading.Lock()
_WS_LOG_ENTRIES: deque[dict[str, Any]] = deque(maxlen=200) _WS_LOG_ENTRIES: deque[dict[str, Any]] = deque(maxlen=200)
_WS_LOG_MAX_RETURN = 100 _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})") _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_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]: def clear_local_alarm_subscription_manual(db: Session) -> dict[str, Any]:
"""Drop persisted/in-memory subscription without calling UME delete.""" """Drop persisted/in-memory subscription without calling UME delete."""
clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY) clear_subscription(db, cache_key=DEFAULT_SUBSCRIPTION_KEY)
@ -579,6 +618,11 @@ def run_alarm_ws_consumer_loop(
time.sleep(1.0) time.sleep(1.0)
continue 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(): if is_ume_subscription_lost():
with _ume_lost_lock: with _ume_lost_lock:
lost_detail = str(_ume_subscription_lost_reason or "") lost_detail = str(_ume_subscription_lost_reason or "")

View file

@ -4,12 +4,17 @@ import json
import logging import logging
import hashlib import hashlib
import re import re
import threading
import time
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any from typing import Any
_sync_log = logging.getLogger("netx.ume.sync") _sync_log = logging.getLogger("netx.ume.sync")
from sqlalchemy import func
from sqlalchemy import text as sql_text 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 sqlalchemy.orm import Session
from .config import settings from .config import settings
@ -196,6 +201,127 @@ def _is_alarm_cleared(alarm: dict[str, Any]) -> bool:
return text in {"true", "1", "yes"} 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: def extract_alarm_from_notification(payload: dict[str, Any]) -> dict[str, Any] | None:
"""Parse alarm-notification from a WS/REST notification envelope.""" """Parse alarm-notification from a WS/REST notification envelope."""
if not isinstance(payload, dict): 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 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. Apply one alarm to ume_alarms_current.
Returns (action, changed) where action is inserted|updated|deleted|skipped. 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 return "skipped", False
existing = db.get(UmeAlarmCurrent, key) existing = db.get(UmeAlarmCurrent, key)
if existing is None: if existing is None:
_mark_alarm_cleared_tombstone(key)
return "deleted", False return "deleted", False
db.delete(existing) db.delete(existing)
_mark_alarm_cleared_tombstone(key)
return "deleted", True return "deleted", True
key = _alarm_key(norm) key = _alarm_key(norm)
existing = db.get(UmeAlarmCurrent, key) if not key:
if existing is None: return "skipped", False
existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=touch_ts) if str(source or "").strip().lower() == "rest" and _is_alarm_cleared_tombstone(key):
db.add(existing) return "skipped", False
action = "inserted"
else: return _upsert_alarm_current(db, key, norm, touch_ts=touch_ts)
action = "updated"
existing.ne_id = _s(_derive_ne_id_from_alarm(norm))
existing.host_name = _lookup_host_name(db, existing.ne_id)
existing.object_name = _s(_pick(norm, "objectName", "object-name"))
existing.event_type = _s(_pick(norm, "eventType", "event-type"))
existing.native_probable_cause = _s(_pick(norm, "nativeProbableCause", "native-probable-cause"))
existing.perceived_severity = _s(_pick(norm, "perceivedSeverity", "perceived-severity"))
existing.is_cleared = _s(_pick(norm, "isCleared", "is-cleared"))
existing.time_created = _s(_pick(norm, "timeCreated", "time-created"))
existing.root_cause_alarm_indication = _s(
_pick(norm, "rootCauseAlarmIndication", "root-cause-alarm-indication")
)
existing.last_seen_at = touch_ts
existing.raw_json = json.dumps(norm, ensure_ascii=False, default=str)
return action, True
def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob:
@ -471,12 +590,34 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = "
return job 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( def _sync_alarms_common(
db: Session, db: Session,
client: UMEClient, client: UMEClient,
*, *,
is_uncleared: bool, is_uncleared: bool,
trigger_mode: str, trigger_mode: str,
wss_active: bool = False,
) -> tuple[UmeSyncJob, UmeAlarmBatch]: ) -> tuple[UmeSyncJob, UmeAlarmBatch]:
domain = "alarms_history" if is_uncleared else "alarms_current" domain = "alarms_history" if is_uncleared else "alarms_current"
job = _build_sync_job(domain, trigger_mode) job = _build_sync_job(domain, trigger_mode)
@ -498,6 +639,8 @@ def _sync_alarms_common(
pulled = inserted = updated = 0 pulled = inserted = updated = 0
deleted_stale_current = 0 deleted_stale_current = 0
host_names_backfilled = 0 host_names_backfilled = 0
reconcile_mode = ""
seen_keys: set[str] = set()
paging_mode = "marker" paging_mode = "marker"
paging_note = "" paging_note = ""
page_no = 0 page_no = 0
@ -506,6 +649,7 @@ def _sync_alarms_common(
graceful_end_by_iterator_error = False graceful_end_by_iterator_error = False
warnings: list[str] = [] warnings: list[str] = []
meta: dict[str, Any] = {} meta: dict[str, Any] = {}
sync_batch_ts = _utc_now_naive()
try: try:
limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000)
limit_max = max(1, limit_max) limit_max = max(1, limit_max)
@ -545,13 +689,22 @@ def _sync_alarms_common(
iterator_500_as_end=iterator_500_as_end, iterator_500_as_end=iterator_500_as_end,
) )
sync_batch_ts = _utc_now_naive() sync_batch_ts = _utc_now_naive()
seen_keys = set()
for rows in pages: for rows in pages:
pulled += len(rows) pulled += len(rows)
for alarm in rows: for alarm in rows:
if is_uncleared: if is_uncleared:
upsert_alarm_history(alarm, touch_ts=sync_batch_ts) upsert_alarm_history(alarm, touch_ts=sync_batch_ts)
else: 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": if action == "inserted":
inserted += 1 inserted += 1
elif action == "updated": elif action == "updated":
@ -564,11 +717,15 @@ def _sync_alarms_common(
paging_note = str(meta.get("paging_note") or "") paging_note = str(meta.get("paging_note") or "")
warnings = [str(x) for x in (meta.get("warnings") or []) if str(x)] 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): if not is_uncleared and _snapshot_reconcile_ok(meta):
deleted_stale_current = int( if wss_active:
db.query(UmeAlarmCurrent) reconcile_mode = "upsert_only"
.filter(UmeAlarmCurrent.last_seen_at < sync_batch_ts) deleted_stale_current = _reconcile_stale_current_alarms(
.delete(synchronize_session=False) db,
sync_batch_ts=sync_batch_ts,
seen_keys=seen_keys,
wss_active=wss_active,
) )
alarm_model = UmeAlarmHistory if is_uncleared else UmeAlarmCurrent alarm_model = UmeAlarmHistory if is_uncleared else UmeAlarmCurrent
@ -593,6 +750,9 @@ def _sync_alarms_common(
"deleted_stale_current_alarms": int(deleted_stale_current), "deleted_stale_current_alarms": int(deleted_stale_current),
"host_names_backfilled": int(host_names_backfilled), "host_names_backfilled": int(host_names_backfilled),
"current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta), "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, ensure_ascii=False,
) )
@ -600,6 +760,7 @@ def _sync_alarms_common(
job.status = "done" job.status = "done"
except Exception as exc: except Exception as exc:
msg = str(exc)[:1024] msg = str(exc)[:1024]
reconcile_mode = "failed"
batch.status = "failed" batch.status = "failed"
batch.error_message = msg batch.error_message = msg
batch.ended_at = _utc_now_naive() batch.ended_at = _utc_now_naive()
@ -625,6 +786,9 @@ def _sync_alarms_common(
"deleted_stale_current_alarms": int(deleted_stale_current), "deleted_stale_current_alarms": int(deleted_stale_current),
"host_names_backfilled": int(host_names_backfilled), "host_names_backfilled": int(host_names_backfilled),
"current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta), "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, ensure_ascii=False,
) )
@ -634,8 +798,24 @@ def _sync_alarms_common(
return job, batch return job, batch
def sync_alarms_current(db: Session, client: UMEClient, *, trigger_mode: str = "manual") -> tuple[UmeSyncJob, UmeAlarmBatch]: def sync_alarms_current(
return _sync_alarms_common(db, client, is_uncleared=False, trigger_mode=trigger_mode) 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( def sync_alarms_history_full(

View file

@ -388,7 +388,7 @@ class UmeSyncServiceTests(unittest.TestCase):
) )
self.db.commit() 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.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-KEEP"))
self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-GONE")) self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-GONE"))
@ -421,7 +421,7 @@ class UmeSyncServiceTests(unittest.TestCase):
) )
self.db.commit() 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-NEW"))
self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-STALE")) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-STALE"))
@ -458,11 +458,11 @@ class UmeSyncServiceTests(unittest.TestCase):
) )
self.db.commit() 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.status, "done")
self.assertEqual(job1.inserted_count, 1) 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.status, "done")
self.assertEqual(job2.updated_count, 1) self.assertEqual(job2.updated_count, 1)
@ -502,7 +502,7 @@ class UmeSyncServiceTests(unittest.TestCase):
svc.settings.ume_page_size = 2 svc.settings.ume_page_size = 2
svc.settings.ume_max_pages = 10 svc.settings.ume_max_pages = 10
try: 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: finally:
svc.settings.ume_page_size = old_page_size svc.settings.ume_page_size = old_page_size
svc.settings.ume_max_pages = old_max_pages 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_marker_max_pages = 10
svc.settings.ume_iterator_500_as_end = True svc.settings.ume_iterator_500_as_end = True
try: 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: finally:
svc.settings.ume_marker_page_limit = old_page_size svc.settings.ume_marker_page_limit = old_page_size
svc.settings.ume_marker_max_pages = old_max_pages svc.settings.ume_marker_max_pages = old_max_pages
@ -568,7 +568,7 @@ class UmeSyncServiceTests(unittest.TestCase):
) )
c = _C() 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.status, "done")
self.assertEqual(job.pulled_count, 1) self.assertEqual(job.pulled_count, 1)
self.assertEqual(c.calls, 1) self.assertEqual(c.calls, 1)
@ -594,7 +594,7 @@ class UmeSyncServiceTests(unittest.TestCase):
) )
c = _C() 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.status, "done")
self.assertEqual(job.pulled_count, 1) self.assertEqual(job.pulled_count, 1)
self.assertEqual(c.calls, 1) self.assertEqual(c.calls, 1)
@ -934,6 +934,119 @@ class UmeAlarmNotificationTests(unittest.TestCase):
self.assertIsNone(load_subscription(db)) self.assertIsNone(load_subscription(db))
db.close() 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): def test_process_alarm_notification_via_ws_helper(self):
from datetime import datetime from datetime import datetime

View file

@ -248,6 +248,11 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
alarmSub.server_subscription_lost_reason ?? 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 = const wsPillLevel =
serverSubLost || wsState === "subscription_lost" serverSubLost || wsState === "subscription_lost"
? "warn" ? "warn"
@ -384,6 +389,11 @@ export function UmePage({ toastOk, toastError }: UmePageProps) {
</span> </span>
) : null} ) : null}
</div> </div>
{scheduledSyncSkipped && !serverSubLost ? (
<div className="pill pill--low" style={{ marginTop: 8 }}>
当前告警由 WSS 实时维护(模式: {currentAlarmsMode}),定时 REST 全量同步已暂停。需要全量对账时请使用下方「同步当前告警」。
</div>
) : null}
{serverSubLost ? ( {serverSubLost ? (
<div className="pill pill--medium" style={{ marginTop: 8 }}> <div className="pill pill--medium" style={{ marginTop: 8 }}>
UME 服务器侧告警订阅已丢失或已过期,本地记录可能仍显示「已建立」。请确认清除本地记录后重新订阅。 UME 服务器侧告警订阅已丢失或已过期,本地记录可能仍显示「已建立」。请确认清除本地记录后重新订阅。

View file

@ -112,6 +112,9 @@ export type UmeAlarmSubscriptionStatus = {
topic?: string; topic?: string;
server_subscription_lost?: boolean; server_subscription_lost?: boolean;
server_subscription_lost_reason?: string; server_subscription_lost_reason?: string;
current_alarms_mode?: "wss" | "rest";
wss_active_for_current_alarms?: boolean;
scheduled_sync_skipped?: boolean;
needs_local_cleanup?: boolean; needs_local_cleanup?: boolean;
ume_already_missing?: boolean; ume_already_missing?: boolean;
cleared_local?: boolean; cleared_local?: boolean;