diff --git a/netx_api/main.py b/netx_api/main.py index 0ac8809..de4562a 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -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] = {} diff --git a/netx_api/mcp_server.py b/netx_api/mcp_server.py index b81d231..e6690de 100644 --- a/netx_api/mcp_server.py +++ b/netx_api/mcp_server.py @@ -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 ""), diff --git a/netx_api/models.py b/netx_api/models.py index 70a0d69..aeb50ce 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -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="") diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 8e4c16c..56ce50f 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -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, diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 2fe532d..5093339 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -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()