feat(ume): denormalize host_name onto alarms for display and grouping

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-20 00:05:10 +08:00
parent 95db646403
commit 21f7cfba41
5 changed files with 136 additions and 6 deletions

View file

@ -360,6 +360,31 @@ def _aggregate_rows(items: list[Any], key_fn) -> list[dict[str, Any]]:
return [{"key": k, "count": v} for k, v in sorted(bucket.items(), key=lambda kv: kv[1], reverse=True)]
def _ume_alarm_host_name(
alarm: UmeAlarmCurrent | UmeAlarmHistory,
ne: UmeInventoryNE | None = None,
) -> str:
hn = str(getattr(alarm, "host_name", "") or "").strip()
if hn:
return hn
if ne is not None:
return str(getattr(ne, "host_name", "") or "").strip()
return ""
def _ume_alarm_ne_group_key(
alarm: UmeAlarmCurrent | UmeAlarmHistory,
ne: UmeInventoryNE | None,
) -> str:
return (
_ume_alarm_host_name(alarm, ne)
or (str(ne.user_label if ne else "") or "").strip()
or (str(ne.ne_name if ne else "") or "").strip()
or str(alarm.ne_id or "").strip()
or "unknown"
)
_PROTOCOL_BUCKET_ZH: dict[str, str] = {
"IP/MPLS": "IP/MPLS",
"ETH": "ETH",
@ -596,6 +621,14 @@ def on_startup() -> None:
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ALTER COLUMN is_cleared TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ALTER COLUMN time_created TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ALTER COLUMN root_cause_alarm_indication TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''")
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_ume_alarms_current_host_name ON ume_alarms_current (host_name)"
)
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_ume_alarms_history_host_name ON ume_alarms_history (host_name)"
)
conn.exec_driver_sql("ALTER TABLE ume_alarms_history ALTER COLUMN alarm_key TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history ALTER COLUMN object_name TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history ALTER COLUMN event_type TYPE TEXT")
@ -1093,13 +1126,16 @@ def ume_list_alarms(
stmt = stmt.filter(UmeAlarmCurrent.ne_id == str(ne_id).strip())
hn = str(host_name or "").strip()
if hn:
stmt = stmt.filter(UmeInventoryNE.host_name.contains(hn))
stmt = stmt.filter(
UmeAlarmCurrent.host_name.contains(hn) | UmeInventoryNE.host_name.contains(hn)
)
kw = str(keyword or "").strip()
if kw:
stmt = stmt.filter(
UmeAlarmCurrent.alarm_key.contains(kw)
| UmeAlarmCurrent.object_name.contains(kw)
| UmeAlarmCurrent.native_probable_cause.contains(kw)
| UmeAlarmCurrent.host_name.contains(kw)
| UmeInventoryNE.ne_name.contains(kw)
| UmeInventoryNE.user_label.contains(kw)
| UmeInventoryNE.ip_address.contains(kw)
@ -1122,7 +1158,7 @@ def ume_list_alarms(
"ne_id": str(alarm.ne_id or ""),
"ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str((ne.user_label if ne else "") or ""),
"host_name": str((ne.host_name if ne else "") or ""),
"host_name": _ume_alarm_host_name(alarm, ne),
"ne_type": str((ne.ne_type if ne else "") or ""),
"object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""),
@ -1210,6 +1246,10 @@ def _extract_ume_raw_group_field(alarm: UmeAlarmCurrent, ne: UmeInventoryNE | No
attr = key[len("ne_") :]
if key == "ne_exists":
return "1" if ne is not None else "0"
if key == "ne_host_name":
hn = str(getattr(alarm, "host_name", "") or "").strip()
if hn:
return hn
if ne is None:
return ""
return str(getattr(ne, attr, "") or "")
@ -1392,7 +1432,7 @@ def ume_alarms_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]:
UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id
).all()
by_severity = _aggregate_rows(rows, lambda x: x[0].perceived_severity)
by_ne = _aggregate_rows(rows, lambda x: (x[1].user_label if x[1] else "") or (x[1].ne_name if x[1] else "") or x[0].ne_id)
by_ne = _aggregate_rows(rows, lambda x: _ume_alarm_ne_group_key(x[0], x[1]))
return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne}
@ -1406,7 +1446,7 @@ def ume_diagnostics(
).all()
by_severity = _aggregate_rows(rows, lambda x: x[0].perceived_severity)
by_alarm_code = _aggregate_rows(rows, lambda x: x[0].event_type)[:10]
by_ne = _aggregate_rows(rows, lambda x: (x[1].user_label if x[1] else "") or (x[1].ne_name if x[1] else "") or x[0].ne_id)[:10]
by_ne = _aggregate_rows(rows, lambda x: _ume_alarm_ne_group_key(x[0], x[1]))[:10]
lang_norm = _normalize_netx_lang(lang)
proto_counts: dict[str, int] = {}

View file

@ -363,9 +363,11 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]:
stmt = stmt.filter(
UmeAlarmCurrent.alarm_key.contains(keyword)
| UmeAlarmCurrent.object_name.contains(keyword)
| UmeAlarmCurrent.host_name.contains(keyword)
| UmeInventoryNE.ne_name.contains(keyword)
| UmeInventoryNE.user_label.contains(keyword)
| UmeInventoryNE.ip_address.contains(keyword)
| UmeInventoryNE.host_name.contains(keyword)
)
page = max(1, int(args.get("page") or 1))
page_size = min(500, max(1, int(args.get("page_size") or 50)))
@ -388,9 +390,12 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]:
{
"alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(alarm.ne_id or ""),
"host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""),
"ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""),
"native_probable_cause": str(alarm.native_probable_cause or ""),
"perceived_severity": str(alarm.perceived_severity or ""),
"is_cleared": str(alarm.is_cleared or ""),
"time_created": str(alarm.time_created or ""),
@ -414,9 +419,11 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]:
stmt = stmt.filter(
UmeAlarmHistory.alarm_key.contains(keyword)
| UmeAlarmHistory.object_name.contains(keyword)
| UmeAlarmHistory.host_name.contains(keyword)
| UmeInventoryNE.ne_name.contains(keyword)
| UmeInventoryNE.user_label.contains(keyword)
| UmeInventoryNE.ip_address.contains(keyword)
| UmeInventoryNE.host_name.contains(keyword)
)
page = max(1, int(args.get("page") or 1))
page_size = min(500, max(1, int(args.get("page_size") or 50)))
@ -430,9 +437,12 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]:
{
"alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(alarm.ne_id or ""),
"host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""),
"ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""),
"native_probable_cause": str(alarm.native_probable_cause or ""),
"perceived_severity": str(alarm.perceived_severity or ""),
"is_cleared": str(alarm.is_cleared or ""),
"time_created": str(alarm.time_created or ""),

View file

@ -168,6 +168,7 @@ class UmeAlarmCurrent(Base):
alarm_key: Mapped[str] = mapped_column(Text, primary_key=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
host_name: Mapped[str] = mapped_column(String(256), default="", index=True, comment="网元主机名(同步时从inventory联表写入)")
object_name: Mapped[str] = mapped_column(Text, default="", index=True)
event_type: Mapped[str] = mapped_column(Text, default="")
native_probable_cause: Mapped[str] = mapped_column(Text, default="")
@ -185,6 +186,7 @@ class UmeAlarmHistory(Base):
alarm_key: Mapped[str] = mapped_column(Text, primary_key=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=True)
host_name: Mapped[str] = mapped_column(String(256), default="", index=True, comment="网元主机名(同步时从inventory联表写入)")
object_name: Mapped[str] = mapped_column(Text, default="", index=True)
event_type: Mapped[str] = mapped_column(Text, default="")
native_probable_cause: Mapped[str] = mapped_column(Text, default="")

View file

@ -9,6 +9,7 @@ from typing import Any
_sync_log = logging.getLogger("netx.ume.sync")
from sqlalchemy import text as sql_text
from sqlalchemy.orm import Session
from .config import settings
@ -42,6 +43,59 @@ def _pick(d: dict[str, Any], *keys: str) -> Any:
return None
def _lookup_host_name(db: Session, ne_id: str) -> str:
nid = _s(ne_id)
if not nid:
return ""
row = db.get(UmeInventoryNE, nid)
if row is None:
return ""
return _s(row.host_name)
def _propagate_host_name_to_alarms(db: Session, ne_id: str, host_name: str) -> None:
nid = _s(ne_id)
if not nid:
return
hn = _s(host_name)
db.query(UmeAlarmCurrent).filter(UmeAlarmCurrent.ne_id == nid).update(
{UmeAlarmCurrent.host_name: hn},
synchronize_session=False,
)
db.query(UmeAlarmHistory).filter(UmeAlarmHistory.ne_id == nid).update(
{UmeAlarmHistory.host_name: hn},
synchronize_session=False,
)
def _backfill_alarm_host_names(db: Session, model: type[UmeAlarmCurrent] | type[UmeAlarmHistory]) -> int:
"""Set alarm.host_name from ume_inventory_ne for all rows with matching ne_id."""
table = str(getattr(model, "__tablename__", "") or "")
if not table:
return 0
bind = db.get_bind()
if bind is not None and str(bind.dialect.name).lower() == "postgresql":
res = db.execute(
sql_text(
f"""
UPDATE {table} AS a
SET host_name = COALESCE(NULLIF(TRIM(ne.host_name), ''), '')
FROM ume_inventory_ne AS ne
WHERE a.ne_id <> '' AND a.ne_id = ne.ne_id
"""
)
)
return int(res.rowcount or 0)
ne_map = {str(r.ne_id or ""): _s(r.host_name) for r in db.query(UmeInventoryNE).all() if str(r.ne_id or "")}
updated = 0
for alarm in db.query(model).filter(model.ne_id != "").all(): # type: ignore[arg-type]
hn = ne_map.get(str(alarm.ne_id or ""), "")
if str(alarm.host_name or "") != hn:
alarm.host_name = hn
updated += 1
return updated
def _alarm_key(alarm: dict[str, Any]) -> str:
key = _s(_pick(alarm, "alarmKey", "alarm-key","alarmkey","id"))
if key:
@ -277,6 +331,7 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = "
existing.vendor = _s(_pick(row, "vendor-name")) or "ZTE"
existing.last_seen_at = now
existing.raw_json = json.dumps(row, ensure_ascii=False, default=str)
_propagate_host_name_to_alarms(db, ne_id, existing.host_name)
db.flush()
deleted_ne = 0
@ -342,6 +397,7 @@ def _sync_alarms_common(
)
pulled = inserted = updated = 0
deleted_stale_current = 0
host_names_backfilled = 0
paging_mode = "marker"
paging_note = ""
page_no = 0
@ -378,6 +434,7 @@ def _sync_alarms_common(
else:
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"))
existing.event_type = _s(_pick(alarm, "eventType", "event-type"))
existing.native_probable_cause = _s(_pick(alarm, "nativeProbableCause", "native-probable-cause"))
@ -416,6 +473,9 @@ def _sync_alarms_common(
.delete(synchronize_session=False)
)
alarm_model = UmeAlarmHistory if is_uncleared else UmeAlarmCurrent
host_names_backfilled = _backfill_alarm_host_names(db, alarm_model)
batch.total_rows = int(pulled)
batch.success_rows = int(inserted + updated)
batch.failed_rows = max(0, int(pulled) - int(inserted + updated))
@ -433,6 +493,7 @@ def _sync_alarms_common(
"graceful_end_by_iterator_error": graceful_end_by_iterator_error,
"warnings": warnings,
"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),
},
ensure_ascii=False,
@ -464,6 +525,7 @@ def _sync_alarms_common(
"graceful_end_by_iterator_error": graceful_end_by_iterator_error,
"warnings": warnings,
"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),
},
ensure_ascii=False,

View file

@ -342,6 +342,16 @@ class UmeSyncServiceTests(unittest.TestCase):
return rows, _D()
self.db.add(
UmeInventoryNE(
ne_id="NE-1",
ne_name="ne1",
user_label="网元1",
host_name="host-ne-1",
)
)
self.db.commit()
job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual")
self.assertEqual(job1.status, "done")
self.assertEqual(job1.inserted_count, 1)
@ -353,6 +363,7 @@ class UmeSyncServiceTests(unittest.TestCase):
row = self.db.get(UmeAlarmCurrent, "AK-1")
self.assertIsNotNone(row)
self.assertEqual(row.perceived_severity, "critical")
self.assertEqual(row.host_name, "host-ne-1")
def test_sync_current_alarms_marker_pagination(self):
class _C:
@ -547,12 +558,17 @@ class UmeSyncServiceTests(unittest.TestCase):
self.assertIn("ne_exists", set(data["selectable_fields"]))
def test_extract_ume_raw_group_field(self):
alarm = UmeAlarmCurrent(alarm_key="AK-X", perceived_severity="major")
ne = UmeInventoryNE(ne_id="NE-X", user_label="site-x")
alarm = UmeAlarmCurrent(alarm_key="AK-X", perceived_severity="major", host_name="host-x")
ne = UmeInventoryNE(ne_id="NE-X", user_label="site-x", host_name="inv-host")
self.assertEqual(_extract_ume_raw_group_field(alarm, ne, "alarm_perceived_severity"), "major")
self.assertEqual(_extract_ume_raw_group_field(alarm, ne, "ne_user_label"), "site-x")
self.assertEqual(_extract_ume_raw_group_field(alarm, ne, "ne_host_name"), "host-x")
self.assertEqual(_extract_ume_raw_group_field(alarm, None, "ne_exists"), "0")
def test_ume_alarms_fields_includes_alarm_host_name(self):
data = ume_alarms_fields()
self.assertIn("alarm_host_name", set(data["selectable_fields"]))
if __name__ == "__main__":
unittest.main()