From fe483f51cc2c5da10333aec83373a9f713708340 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 6 Aug 2026 18:17:45 +0800 Subject: [PATCH] Add UME TopoNodes/TopologicalLinks daily sync into local tables. Phase-1 docks both topology APIs without writing Fabric; expose a 24h topology_auto_sync task and domain=topology manual sync for later LLDP alignment. Co-authored-by: Cursor --- .env.example | 5 + netx_api/config.py | 5 + netx_api/mcp/db_server.py | 16 +- netx_api/models/__init__.py | 4 + netx_api/models/ume.py | 42 +++++ netx_api/runtime_task_messages.py | 1 + netx_api/schema_patches.py | 42 +++++ netx_api/ume_client.py | 48 +++++- netx_api/ume_runtime.py | 74 ++++++++- netx_api/ume_support.py | 10 +- netx_api/ume_sync_router.py | 26 ++- netx_api/ume_sync_service.py | 4 +- netx_api/ume_sync_topology.py | 263 ++++++++++++++++++++++++++++++ tests/test_ume_sync.py | 158 +++++++++++++++++- web/src/i18n/en.ts | 2 + web/src/i18n/zh.ts | 2 + web/src/utils/runtimeMessages.ts | 1 + 17 files changed, 690 insertions(+), 13 deletions(-) create mode 100644 netx_api/ume_sync_topology.py diff --git a/.env.example b/.env.example index 6df2c13..37764a8 100644 --- a/.env.example +++ b/.env.example @@ -28,6 +28,11 @@ NETX_UME_TOKEN_HANDSHAKE_PATH=/restconf/operations/zte-security:oauth_handshake NETX_UME_TOKEN_LOGOUT_PATH=/restconf/operations/zte-security:oauth_token NETX_UME_NE_PATH=/restconf/data/zte-resources-module:network-elements NETX_UME_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list +NETX_UME_TOPO_NODES_PATH=/restconf/data/zte-resources-module:TopoNodes +NETX_UME_TOPOLOGICAL_LINKS_PATH=/restconf/data/zte-resources-module:TopologicalLinks +NETX_UME_SYNC_TOPOLOGY_AUTO_ENABLED=true +NETX_UME_SYNC_TOPOLOGY_EVERY_HOURS=24 +# NETX_UME_TOPOLOGY_TIMEOUT_S=120 NETX_UME_SYNC_ALARMS_CURRENT_ENABLED=true NETX_UME_SYNC_ALARMS_CURRENT_INTERVAL_S=18000 NETX_UME_SYNC_ALARMS_CURRENT_SKIP_WHEN_WS=true diff --git a/netx_api/config.py b/netx_api/config.py index a7ba304..07a81a9 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -56,11 +56,16 @@ class Settings(BaseSettings): ume_notification_topic: str = "ALARM" ume_sync_inventory_auto_enabled: bool = True ume_sync_inventory_every_hours: int = 48 + ume_sync_topology_auto_enabled: bool = True + ume_sync_topology_every_hours: int = 24 + ume_topology_timeout_s: float = 120.0 ume_token_path: str = "/restconf/operations/zte-security:oauth_token" ume_token_handshake_path: str = "/restconf/operations/zte-security:oauth_handshake" ume_token_logout_path: str = "/restconf/operations/zte-security:oauth_token" ume_ne_path: str = "/restconf/data/zte-resources-module:network-elements" ume_alarms_path: str = "/restconf/data/zte-alarms:alarms/alarm-list" + ume_topo_nodes_path: str = "/restconf/data/zte-resources-module:TopoNodes" + ume_topological_links_path: str = "/restconf/data/zte-resources-module:TopologicalLinks" ume_sync_alarms_history_every_hours: int = 24 # Managed NE credentials (Fernet key; generate with cryptography.fernet.Fernet.generate_key()) credential_secret_key: str = "" diff --git a/netx_api/mcp/db_server.py b/netx_api/mcp/db_server.py index 71bf228..7e1b8f4 100644 --- a/netx_api/mcp/db_server.py +++ b/netx_api/mcp/db_server.py @@ -8,7 +8,7 @@ from .db import Base, SessionLocal, engine from .importer import aggregate_alarms, query_alarms from .models import AlarmBatch, UmeAlarmCurrent, UmeAlarmHistory, UmeInventoryNE from .ume_client import UMEClient -from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full +from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full, sync_topology_full from .ume_token_store import ( clear_shared_token, load_shared_token, @@ -103,7 +103,7 @@ def _tool_list() -> list[dict[str, Any]]: "properties": { "domains": { "type": "array", - "items": {"type": "string", "enum": ["inventory", "alarms_current", "alarms_history"]}, + "items": {"type": "string", "enum": ["inventory", "alarms_current", "alarms_history", "topology"]}, }, "trigger_mode": {"type": "string", "enum": ["manual", "schedule"]}, }, @@ -299,6 +299,18 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: "error_message": str(j.error_message or ""), } ) + if "topology" in domains: + j = sync_topology_full(db, client, trigger_mode=trigger_mode) + results.append( + { + "domain": "topology", + "status": j.status, + "pulled_count": int(j.pulled_count or 0), + "inserted_count": int(j.inserted_count or 0), + "updated_count": int(j.updated_count or 0), + "error_message": str(j.error_message or ""), + } + ) return {"content": [{"type": "text", "text": json.dumps({"ok": True, "jobs": results}, ensure_ascii=False)}]} if name == "umeListNE": stmt = db.query(UmeInventoryNE) diff --git a/netx_api/models/__init__.py b/netx_api/models/__init__.py index 3f8f9c4..2d5e34b 100644 --- a/netx_api/models/__init__.py +++ b/netx_api/models/__init__.py @@ -56,6 +56,8 @@ from .ume import ( UmeKeyAlertRule, UmeSyncJob, UmeTokenCache, + UmeTopoLink, + UmeTopoNode, ) __all__ = [ @@ -74,6 +76,8 @@ __all__ = [ "UmeKeyAlertForwardLog", "UmeAlarmSubscription", "UmeTokenCache", + "UmeTopoNode", + "UmeTopoLink", "ManagedNE", "CliConnectProfile", "UmeCliOverride", diff --git a/netx_api/models/ume.py b/netx_api/models/ume.py index e0723a9..7a920a6 100644 --- a/netx_api/models/ume.py +++ b/netx_api/models/ume.py @@ -171,3 +171,45 @@ class UmeTokenCache(Base): lock_owner: Mapped[str] = mapped_column(String(128), default="", index=True) lock_expires_at_epoch_s: Mapped[int] = mapped_column(Integer, default=0, index=True) updated_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + + +class UmeTopoNode(Base): + """UME TopoNodes snapshot (coordinates / subnet tree). Phase-1 dock only.""" + + __tablename__ = "ume_topo_node" + + node_id: Mapped[str] = mapped_column(String(128), primary_key=True, comment="UME nodeId") + name: Mapped[str] = mapped_column(String(512), default="", index=True) + node_type: Mapped[str] = mapped_column(String(64), default="", index=True, comment="TOPO_NODE_ME|TOPO_NODE_SBN") + user_label: Mapped[str] = mapped_column(String(512), default="", index=True) + owner: Mapped[str] = mapped_column(String(64), default="") + parent_node: Mapped[str] = mapped_column(String(512), default="", index=True) + x_pos: Mapped[int | None] = mapped_column(Integer, nullable=True) + y_pos: Mapped[int | None] = mapped_column(Integer, nullable=True) + ume_ne_id: Mapped[str] = mapped_column(String(128), default="", index=True, comment="ME uuid when nodeType=ME") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeTopoLink(Base): + """UME TopologicalLinks snapshot. Phase-1 dock only (no Fabric apply yet).""" + + __tablename__ = "ume_topo_link" + + link_id: Mapped[str] = mapped_column(String(128), primary_key=True, comment="UME linkId") + name: Mapped[str] = mapped_column(String(1024), default="", index=True) + user_label: Mapped[str] = mapped_column(Text, default="") + owner: Mapped[str] = mapped_column(String(64), default="") + direction: Mapped[str] = mapped_column(String(32), default="") + layer_rate: Mapped[int | None] = mapped_column(Integer, nullable=True) + connection_status: Mapped[str] = mapped_column(String(64), default="", index=True) + a_end_tp_ref: Mapped[str] = mapped_column(Text, default="") + z_end_tp_ref: Mapped[str] = mapped_column(Text, default="") + a_ume_ne_id: Mapped[str] = mapped_column(String(128), default="", index=True) + z_ume_ne_id: Mapped[str] = mapped_column(String(128), default="", index=True) + a_ptp: Mapped[str] = mapped_column(String(256), default="") + z_ptp: Mapped[str] = mapped_column(String(256), default="") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=utcnow_naive, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") diff --git a/netx_api/runtime_task_messages.py b/netx_api/runtime_task_messages.py index 5e76152..14d2df5 100644 --- a/netx_api/runtime_task_messages.py +++ b/netx_api/runtime_task_messages.py @@ -6,6 +6,7 @@ RT_WSS_ACTIVE_SKIP_REST = "rt:wss_active_skip_rest" RT_PULLING_ALARMS_CURRENT = "rt:pulling_alarms_current" RT_ALARMS_SYNC_IN_PROGRESS_SKIP = "rt:alarms_sync_in_progress_skip" RT_PULLING_INVENTORY = "rt:pulling_inventory" +RT_PULLING_TOPOLOGY = "rt:pulling_topology" RT_UME_WS_DISABLED_NO_BASE_URL = "rt:ume_ws_disabled_no_base_url" RT_OCLAW_FWD_DISABLED = "rt:oclaw_fwd_disabled" RT_RESUMED_SYNC_SOON = "rt:resumed_sync_soon" diff --git a/netx_api/schema_patches.py b/netx_api/schema_patches.py index a17d2e0..7e758de 100644 --- a/netx_api/schema_patches.py +++ b/netx_api/schema_patches.py @@ -344,6 +344,48 @@ def apply_domain_schema_patches(conn: Connection) -> None: ) """, "CREATE INDEX IF NOT EXISTS ix_ume_cli_override_connect_status ON ume_cli_override (connect_status)", + """ + CREATE TABLE IF NOT EXISTS ume_topo_node ( + node_id VARCHAR(128) PRIMARY KEY, + name VARCHAR(512) DEFAULT '', + node_type VARCHAR(64) DEFAULT '', + user_label VARCHAR(512) DEFAULT '', + owner VARCHAR(64) DEFAULT '', + parent_node VARCHAR(512) DEFAULT '', + x_pos INTEGER, + y_pos INTEGER, + ume_ne_id VARCHAR(128) DEFAULT '', + first_seen_at TIMESTAMP, + last_seen_at TIMESTAMP, + raw_json TEXT DEFAULT '{}' + ) + """, + "CREATE INDEX IF NOT EXISTS ix_ume_topo_node_node_type ON ume_topo_node (node_type)", + "CREATE INDEX IF NOT EXISTS ix_ume_topo_node_ume_ne_id ON ume_topo_node (ume_ne_id)", + "CREATE INDEX IF NOT EXISTS ix_ume_topo_node_last_seen_at ON ume_topo_node (last_seen_at)", + """ + CREATE TABLE IF NOT EXISTS ume_topo_link ( + link_id VARCHAR(128) PRIMARY KEY, + name VARCHAR(1024) DEFAULT '', + user_label TEXT DEFAULT '', + owner VARCHAR(64) DEFAULT '', + direction VARCHAR(32) DEFAULT '', + layer_rate INTEGER, + connection_status VARCHAR(64) DEFAULT '', + a_end_tp_ref TEXT DEFAULT '', + z_end_tp_ref TEXT DEFAULT '', + a_ume_ne_id VARCHAR(128) DEFAULT '', + z_ume_ne_id VARCHAR(128) DEFAULT '', + a_ptp VARCHAR(256) DEFAULT '', + z_ptp VARCHAR(256) DEFAULT '', + first_seen_at TIMESTAMP, + last_seen_at TIMESTAMP, + raw_json TEXT DEFAULT '{}' + ) + """, + "CREATE INDEX IF NOT EXISTS ix_ume_topo_link_a_ume_ne_id ON ume_topo_link (a_ume_ne_id)", + "CREATE INDEX IF NOT EXISTS ix_ume_topo_link_z_ume_ne_id ON ume_topo_link (z_ume_ne_id)", + "CREATE INDEX IF NOT EXISTS ix_ume_topo_link_last_seen_at ON ume_topo_link (last_seen_at)", "COMMENT ON TABLE ume_inventory_ne IS '网元对象详细信息'", "COMMENT ON COLUMN ume_inventory_ne.ne_id IS '网元uuid'", "COMMENT ON COLUMN ume_inventory_ne.ne_name IS '资源名称'", diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index 30dfbf7..4116ffd 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -96,6 +96,12 @@ class UMEClient: self.token_logout_path = str(token_logout_path if token_logout_path is not None else settings.ume_token_logout_path).strip() self.ne_path = str(ne_path if ne_path is not None else settings.ume_ne_path).strip() self.alarms_path = str(alarms_path if alarms_path is not None else settings.ume_alarms_path).strip() + self.topo_nodes_path = str(getattr(settings, "ume_topo_nodes_path", "") or "").strip() or ( + "/restconf/data/zte-resources-module:TopoNodes" + ) + self.topological_links_path = str(getattr(settings, "ume_topological_links_path", "") or "").strip() or ( + "/restconf/data/zte-resources-module:TopologicalLinks" + ) self.notification_establish_path = str( notification_establish_path if notification_establish_path is not None @@ -231,11 +237,12 @@ class UMEClient: headers[self.auth_header] = token return headers - def _client(self) -> httpx.Client: + def _client(self, *, timeout_s: float | None = None) -> httpx.Client: # Use explicit HTTPTransport to keep behavior consistent with onsite validation. # In this mode, requests run over HTTP/1.1 and avoid HTTP/2 negotiation issues. transport = httpx.HTTPTransport(verify=self.verify_tls, http2=False) - return httpx.Client(transport=transport, timeout=self.timeout_s) + to = self.timeout_s if timeout_s is None else max(3.0, float(timeout_s)) + return httpx.Client(transport=transport, timeout=to) def _extract_token_and_ttl(self, payload: dict[str, Any]) -> tuple[str, int | None]: @@ -462,6 +469,7 @@ class UMEClient: *, params: dict[str, Any] | None = None, body: dict[str, Any] | None = None, + timeout_s: float | None = None, ) -> tuple[dict[str, Any], RequestDiagnostics]: self.refresh_if_needed() url = self._build_url(path) @@ -469,12 +477,12 @@ class UMEClient: retry_count = 0 t0 = time() try: - with self._client() as client: + with self._client(timeout_s=timeout_s) as client: resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True)) if resp.status_code in (401, 403): retry_count = 1 self.login(force=True) - with self._client() as client: + with self._client(timeout_s=timeout_s) as client: resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True)) marker = str(resp.headers.get("marker") or "").strip() is_end_raw = str(resp.headers.get("is-end-of-reply") or "").strip().lower() @@ -604,6 +612,38 @@ class UMEClient: rows = self._extract_named_list(data, ["alarm-list", "alarm"]) return rows, diag + def _topology_timeout_s(self) -> float: + base = float(getattr(settings, "ume_topology_timeout_s", 0) or 0) + if base > 0: + return max(3.0, base) + return max(self.timeout_s, 60.0) + + def get_topo_nodes(self) -> tuple[list[dict[str, Any]], RequestDiagnostics]: + data, diag = self.request_json("GET", self.topo_nodes_path, timeout_s=self._topology_timeout_s()) + rows = self._extract_named_list(data, ["TopoNodes", "TopoNode", "topo-nodes", "topo-node"]) + if rows: + return rows, diag + for v in data.values(): + lst = _coerce_list(v) + if lst and isinstance(lst[0], dict): + return [x for x in lst if isinstance(x, dict)], diag + return [], diag + + def get_topological_links(self) -> tuple[list[dict[str, Any]], RequestDiagnostics]: + data, diag = self.request_json( + "GET", self.topological_links_path, timeout_s=self._topology_timeout_s() + ) + rows = self._extract_named_list( + data, ["TopologicalLinks", "TopologicalLink", "topological-links", "topological-link"] + ) + if rows: + return rows, diag + for v in data.values(): + lst = _coerce_list(v) + if lst and isinstance(lst[0], dict): + return [x for x in lst if isinstance(x, dict)], diag + return [], diag + def _extract_subscription_output(self, payload: dict[str, Any]) -> tuple[str, str]: sub_id = "" uri = "" diff --git a/netx_api/ume_runtime.py b/netx_api/ume_runtime.py index 6eb46bd..3625fa0 100644 --- a/netx_api/ume_runtime.py +++ b/netx_api/ume_runtime.py @@ -24,13 +24,14 @@ from .ume_alarm_ws import ( load_persisted_subscription, start_ume_alarm_ws_consumer, ) -from .ume_sync_service import sync_alarms_current, sync_inventory_full +from .ume_sync_service import sync_alarms_current, sync_inventory_full, sync_topology_full from .runtime_task_messages import ( RT_ALARMS_SYNC_IN_PROGRESS_SKIP, RT_KEEPALIVE_FAILED, RT_OCLAW_FWD_DISABLED, RT_PULLING_ALARMS_CURRENT, RT_PULLING_INVENTORY, + RT_PULLING_TOPOLOGY, RT_STARTUP_GATE_WAITING, RT_UME_WS_DISABLED_NO_BASE_URL, RT_WSS_ACTIVE_SKIP_REST, @@ -286,6 +287,77 @@ def start_api_sideband_threads() -> None: last_run_at=datetime.now(timezone.utc), last_error=f"startup_thread_init_failed: {str(exc)[:180]}", ) + try: + if bool(getattr(settings, "ume_sync_topology_auto_enabled", True)): + hours = int(getattr(settings, "ume_sync_topology_every_hours", 24) or 24) + hours = max(1, min(hours, 168)) + topology_interval_s = int(hours * 3600) + ume_support._refresh_runtime_task_idle("topology_auto_sync", "topology") + + def _topology_auto_sync_loop() -> None: + ume_support._refresh_runtime_task_idle("topology_auto_sync", "topology") + while True: + try: + _schedule_log.info( + "topology_auto_sync: loop tick paused=%s", + ume_support._runtime_is_paused("topology_auto_sync"), + ) + if ume_support._runtime_is_paused("topology_auto_sync"): + time.sleep(1) + continue + ume_support._maybe_wait_for_sync_interval( + task_id="topology_auto_sync", + domain="topology", + interval_s=topology_interval_s, + label="topology_auto_sync", + ) + _schedule_log.info( + "topology_auto_sync: iteration start (interval=%ss)", + topology_interval_s, + ) + ume_support._set_runtime_task( + "topology_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error=RT_PULLING_TOPOLOGY, + ) + db = SessionLocal() + try: + client = ume_support._ume_client() + sync_topology_full(db, client, trigger_mode="schedule") + _schedule_log.info("topology_auto_sync: sync finished ok") + ume_support._set_runtime_task( + "topology_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error="", + ) + finally: + db.close() + except Exception as exc: + _schedule_log.exception("topology_auto_sync: sync failed: %s", exc) + ume_support._set_runtime_task( + "topology_auto_sync", + status="error", + last_run_at=datetime.now(timezone.utc), + last_error=str(exc)[:240], + ) + + t_topo = threading.Thread( + target=_topology_auto_sync_loop, name="ume-topology-auto-sync", daemon=True + ) + t_topo.start() + _schedule_log.info("started thread %s alive=%s", t_topo.name, t_topo.is_alive()) + if not t_topo.is_alive(): + _schedule_log.error("ume-topology-auto-sync thread exited immediately (check uncaught errors above)") + except Exception as exc: + _schedule_log.exception("startup: topology_auto_sync thread init failed: %s", exc) + ume_support._set_runtime_task( + "topology_auto_sync", + status="error", + last_run_at=datetime.now(timezone.utc), + last_error=f"startup_thread_init_failed: {str(exc)[:180]}", + ) try: if bool(getattr(settings, "ume_alarm_ws_enabled", True)) and str(getattr(settings, "ume_base_url", "") or "").strip(): if load_persisted_subscription(): diff --git a/netx_api/ume_support.py b/netx_api/ume_support.py index 68a8136..5e5f71c 100644 --- a/netx_api/ume_support.py +++ b/netx_api/ume_support.py @@ -38,7 +38,7 @@ from .ume_alarm_ws import ( is_wss_active_for_current_alarms, ) from .ume_client import UMEClient -from .ume_sync_service import sync_alarms_current, sync_inventory_full +from .ume_sync_service import sync_alarms_current, sync_inventory_full, sync_topology_full from .ume_token_store import ( clear_shared_token, load_shared_token, @@ -66,6 +66,7 @@ _UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = { "alarms_current_ws_consumer": {"task": "alarms_current_ws_consumer", "status": "init", "last_run_at": None, "last_error": ""}, "oclaw_alarm_forwarder": {"task": "oclaw_alarm_forwarder", "status": "init", "last_run_at": None, "last_error": ""}, "inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, + "topology_auto_sync": {"task": "topology_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, } _UME_WS_STOP_EVENT: threading.Event | None = None _UME_RUNTIME_PAUSED: dict[str, bool] = {} @@ -185,6 +186,13 @@ def _runtime_task_interval_fields(task_id: str) -> tuple[int | None, str]: hours = max(1, min(hours, 168)) eff = int(hours * 3600) return eff, _format_runtime_interval_label(eff) + if task_id == "topology_auto_sync": + if not bool(getattr(settings, "ume_sync_topology_auto_enabled", True)): + return None, "disabled" + hours = int(getattr(settings, "ume_sync_topology_every_hours", 24) or 24) + hours = max(1, min(hours, 168)) + eff = int(hours * 3600) + return eff, _format_runtime_interval_label(eff) return None, "—" diff --git a/netx_api/ume_sync_router.py b/netx_api/ume_sync_router.py index dc8d189..b30be78 100644 --- a/netx_api/ume_sync_router.py +++ b/netx_api/ume_sync_router.py @@ -34,7 +34,7 @@ from .ume_support import ( _set_runtime_task, _ume_client, ) -from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full +from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full, sync_topology_full _log = logging.getLogger("netx.ume.router") router = APIRouter(tags=["ume"]) @@ -65,6 +65,18 @@ def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db "error_message": str(job.error_message or ""), } ) + if "topology" in domain_set: + job = sync_topology_full(db, client, trigger_mode=trigger_mode) + out["jobs"].append( + { + "domain": "topology", + "status": job.status, + "pulled_count": int(job.pulled_count or 0), + "inserted_count": int(job.inserted_count or 0), + "updated_count": int(job.updated_count or 0), + "error_message": str(job.error_message or ""), + } + ) if "alarms" in domain_set or "alarms_current" in domain_set: paused_ws_for_sync = False if is_wss_active_for_current_alarms() and trigger_mode == "manual": @@ -123,7 +135,7 @@ def _ume_sync_job_deleted_count(row: UmeSyncJob) -> int: return 0 if not isinstance(obj, dict): return 0 - inv = cur = 0 + inv = cur = topo_n = topo_l = 0 try: inv = max(0, int(obj.get("deleted_inventory_ne") or 0)) except Exception: @@ -132,7 +144,15 @@ def _ume_sync_job_deleted_count(row: UmeSyncJob) -> int: cur = max(0, int(obj.get("deleted_stale_current_alarms") or 0)) except Exception: pass - return int(inv + cur) + try: + topo_n = max(0, int(obj.get("deleted_topo_nodes") or 0)) + except Exception: + pass + try: + topo_l = max(0, int(obj.get("deleted_topo_links") or 0)) + except Exception: + pass + return int(inv + cur + topo_n + topo_l) @router.get("/v1/ume/sync/status") diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 326c9e1..8908bd9 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -1,4 +1,4 @@ -"""UME sync service facade (inventory + alarms).""" +"""UME sync service facade (inventory + alarms + topology).""" from __future__ import annotations from .config import settings @@ -13,6 +13,7 @@ from .ume_alarm_apply import ( ) from .ume_sync_common import _pick, _s, _utc_now_naive from .ume_sync_pull import sync_alarms_current, sync_alarms_history_full, sync_inventory_full +from .ume_sync_topology import sync_topology_full __all__ = [ "_alarm_key", @@ -29,4 +30,5 @@ __all__ = [ "sync_alarms_current", "sync_alarms_history_full", "sync_inventory_full", + "sync_topology_full", ] diff --git a/netx_api/ume_sync_topology.py b/netx_api/ume_sync_topology.py new file mode 100644 index 0000000..7c25453 --- /dev/null +++ b/netx_api/ume_sync_topology.py @@ -0,0 +1,263 @@ +"""UME TopoNodes + TopologicalLinks sync (phase-1: local tables only).""" +from __future__ import annotations + +import hashlib +import json +import logging +import re +from typing import Any + +from sqlalchemy.orm import Session + +from .models import UmeSyncJob, UmeTopoLink, UmeTopoNode +from .ume_client import UMEClient +from .ume_raw import dumps_ume_raw +from .ume_sync_common import _pick, _s, _utc_now_naive +from .ume_sync_pull import _build_sync_job + +_sync_log = logging.getLogger("netx.ume.sync") + +_UUID_RE = re.compile( + r"[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}" +) +_ME_BRACE_RE = re.compile( + r"ME\{([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})\}", + re.IGNORECASE, +) +_PTP_RE = re.compile(r"PTP=\{([^}]*)\}", re.IGNORECASE) + + +def extract_me_uuid(text: str) -> str: + """Extract managed-element uuid from TP ref or TOPO_NODE_ME* name.""" + s = str(text or "").strip() + if not s: + return "" + m = _ME_BRACE_RE.search(s) + if m: + return m.group(1) + # TOPO_NODE_ME (no braces) or trailing uuid + low = s.upper() + if "TOPO_NODE_ME" in low: + idx = low.find("TOPO_NODE_ME") + rest = s[idx + len("TOPO_NODE_ME") :] + m2 = _UUID_RE.search(rest) + if m2: + return m2.group(0) + m3 = _UUID_RE.search(s) + return m3.group(0) if m3 else "" + + +def extract_ptp(text: str) -> str: + s = str(text or "").strip() + if not s: + return "" + m = _PTP_RE.search(s) + return (m.group(1) if m else "")[:256] + + +def _first_tp_ref(value: Any) -> str: + if isinstance(value, list): + for item in value: + s = _s(item) + if s: + return s + return "" + return _s(value) + + +def _as_optional_int(value: Any) -> int | None: + if value is None or value == "": + return None + try: + return int(value) + except Exception: + return None + + +def _link_id_from_row(row: dict[str, Any]) -> str: + lid = _s(_pick(row, "linkId", "link-id", "link_id", "id")) + if lid: + return lid[:128] + name = _s(_pick(row, "name")) + if not name: + return "" + digest = hashlib.sha1(name.encode("utf-8", errors="replace")).hexdigest()[:40] + return f"name:{digest}"[:128] + + +def _node_id_from_row(row: dict[str, Any]) -> str: + nid = _s(_pick(row, "nodeId", "node-id", "node_id", "id")) + if nid: + return nid[:128] + name = _s(_pick(row, "name")) + if not name: + return "" + digest = hashlib.sha1(name.encode("utf-8", errors="replace")).hexdigest()[:40] + return f"name:{digest}"[:128] + + +def _ume_ne_id_for_topo_node(*, node_type: str, name: str) -> str: + nt = str(node_type or "").strip().upper() + if nt != "TOPO_NODE_ME": + return "" + return extract_me_uuid(name) + + +def sync_topology_full(db: Session, client: UMEClient, *, trigger_mode: str = "manual") -> UmeSyncJob: + """Pull TopoNodes + TopologicalLinks into local tables (no Fabric apply).""" + job = _build_sync_job("topology", trigger_mode) + db.add(job) + db.flush() + db.commit() + _sync_log.info("topology sync job %s committed as running (trigger=%s)", getattr(job, "id", "?"), trigger_mode) + + pulled = inserted = updated = 0 + nodes_pulled = nodes_ins = nodes_upd = nodes_del = 0 + links_pulled = links_ins = links_upd = links_del = 0 + nodes_skip = links_skip = 0 + + try: + now = _utc_now_naive() + + node_rows, node_diag = client.get_topo_nodes() + nodes_pulled = len(node_rows) + seen_nodes: set[str] = set() + for row in node_rows: + if not isinstance(row, dict): + nodes_skip += 1 + continue + node_id = _node_id_from_row(row) + if not node_id: + nodes_skip += 1 + continue + seen_nodes.add(node_id) + name = _s(_pick(row, "name")) + node_type = _s(_pick(row, "nodeType", "node-type", "node_type")) + existing = db.get(UmeTopoNode, node_id) + if existing is None: + existing = UmeTopoNode(node_id=node_id, first_seen_at=now) + db.add(existing) + nodes_ins += 1 + else: + nodes_upd += 1 + existing.name = name[:512] + existing.node_type = node_type[:64] + existing.user_label = _s(_pick(row, "userLabel", "user-label", "user_label"))[:512] + existing.owner = _s(_pick(row, "owner"))[:64] + existing.parent_node = _s(_pick(row, "parentNode", "parent-node", "parent_node"))[:512] + existing.x_pos = _as_optional_int(_pick(row, "xPos", "x-pos", "x_pos")) + existing.y_pos = _as_optional_int(_pick(row, "yPos", "y-pos", "y_pos")) + existing.ume_ne_id = _ume_ne_id_for_topo_node(node_type=node_type, name=name)[:128] + existing.last_seen_at = now + existing.raw_json = dumps_ume_raw(row) + + db.flush() + if seen_nodes: + nodes_del = int( + db.query(UmeTopoNode) + .filter(~UmeTopoNode.node_id.in_(list(seen_nodes))) + .delete(synchronize_session=False) + ) + else: + # Successful empty snapshot → clear local table. + nodes_del = int(db.query(UmeTopoNode).delete(synchronize_session=False)) + + link_rows, link_diag = client.get_topological_links() + links_pulled = len(link_rows) + seen_links: set[str] = set() + for row in link_rows: + if not isinstance(row, dict): + links_skip += 1 + continue + link_id = _link_id_from_row(row) + if not link_id: + links_skip += 1 + continue + seen_links.add(link_id) + a_ref = _first_tp_ref(_pick(row, "aEndTpRefList", "a-end-tp-ref-list", "aEndTpRef")) + z_ref = _first_tp_ref(_pick(row, "zEndTpRefList", "z-end-tp-ref-list", "zEndTpRef")) + existing = db.get(UmeTopoLink, link_id) + if existing is None: + existing = UmeTopoLink(link_id=link_id, first_seen_at=now) + db.add(existing) + links_ins += 1 + else: + links_upd += 1 + existing.name = _s(_pick(row, "name"))[:1024] + existing.user_label = _s(_pick(row, "userLabel", "user-label", "user_label")) + existing.owner = _s(_pick(row, "owner"))[:64] + existing.direction = _s(_pick(row, "direction"))[:32] + existing.layer_rate = _as_optional_int(_pick(row, "layerRate", "layer-rate", "layer_rate")) + existing.connection_status = _s( + _pick(row, "connection-status", "connectionStatus", "connection_status") + )[:64] + existing.a_end_tp_ref = a_ref + existing.z_end_tp_ref = z_ref + existing.a_ume_ne_id = extract_me_uuid(a_ref)[:128] + existing.z_ume_ne_id = extract_me_uuid(z_ref)[:128] + existing.a_ptp = extract_ptp(a_ref) + existing.z_ptp = extract_ptp(z_ref) + existing.last_seen_at = now + existing.raw_json = dumps_ume_raw(row) + + db.flush() + if seen_links: + links_del = int( + db.query(UmeTopoLink) + .filter(~UmeTopoLink.link_id.in_(list(seen_links))) + .delete(synchronize_session=False) + ) + else: + links_del = int(db.query(UmeTopoLink).delete(synchronize_session=False)) + + pulled = nodes_pulled + links_pulled + inserted = nodes_ins + links_ins + updated = nodes_upd + links_upd + + job.details_json = json.dumps( + { + "nodes": { + "pulled": nodes_pulled, + "inserted": nodes_ins, + "updated": nodes_upd, + "deleted": nodes_del, + "skipped": nodes_skip, + "latency_ms": int(getattr(node_diag, "latency_ms", 0) or 0), + }, + "links": { + "pulled": links_pulled, + "inserted": links_ins, + "updated": links_upd, + "deleted": links_del, + "skipped": links_skip, + "latency_ms": int(getattr(link_diag, "latency_ms", 0) or 0), + }, + "deleted_topo_nodes": nodes_del, + "deleted_topo_links": links_del, + }, + ensure_ascii=False, + ) + job.status = "done" + _sync_log.info( + "topology sync done nodes=%s/%s/%s links=%s/%s/%s deleted_n=%s deleted_l=%s", + nodes_pulled, + nodes_ins, + nodes_upd, + links_pulled, + links_ins, + links_upd, + nodes_del, + links_del, + ) + except Exception as exc: + job.status = "failed" + job.error_message = str(exc)[:1024] + _sync_log.exception("topology sync failed: %s", exc) + finally: + job.pulled_count = int(pulled) + job.inserted_count = int(inserted) + job.updated_count = int(updated) + job.ended_at = _utc_now_naive() + db.commit() + db.refresh(job) + return job diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index b5283e6..be40c7f 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -10,7 +10,7 @@ from sqlalchemy.orm import sessionmaker from netx_api.db import Base from netx_api.main import _extract_ume_raw_group_field, _serialize_ume_alarm_raw_row, sql_ume_query, ume_alarms_fields -from netx_api.models import UmeAlarmCurrent, UmeInventoryNE +from netx_api.models import UmeAlarmCurrent, UmeInventoryNE, UmeTopoLink, UmeTopoNode from netx_api.ume_client import UMEClient, _parse_json_response from netx_api import ume_alarm_ws from netx_api.models import UmeAlarmSubscription @@ -36,7 +36,9 @@ from netx_api.ume_sync_service import ( normalize_yang_alarm, sync_alarms_current, sync_inventory_full, + sync_topology_full, ) +from netx_api.ume_sync_topology import extract_me_uuid, extract_ptp from fastapi import HTTPException @@ -1088,5 +1090,159 @@ class UmeAlarmNotificationTests(unittest.TestCase): self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-WS-3")) +class UmeTopologySyncTests(unittest.TestCase): + def setUp(self): + 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) + self.db = TestingSessionLocal() + + def tearDown(self): + self.db.close() + + def test_extract_me_and_ptp(self): + tp = "ME{4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2},EQ={/r=0/sh=1/sl=1},PTP={/p=1_16}" + self.assertEqual(extract_me_uuid(tp), "4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2") + self.assertEqual(extract_ptp(tp), "/p=1_16") + self.assertEqual( + extract_me_uuid("TOPO_NODE_ME7e8ac1c7-9d34-42d3-adfb-b1031e7c145a"), + "7e8ac1c7-9d34-42d3-adfb-b1031e7c145a", + ) + + def test_sync_topology_upsert_and_reconcile(self): + class _Diag: + latency_ms = 1 + + class _CWide: + def get_topo_nodes(self): + return ( + [ + { + "nodeId": "66e807c0-94b6-4c50-a02e-fcb56d08bdea", + "name": "SBN{66e807c0-94b6-4c50-a02e-fcb56d08bdea}", + "nodeType": "TOPO_NODE_SBN", + "owner": "ZTE", + "userLabel": "SC IPRAN Network", + "yPos": 126, + "xPos": 460, + "parentNode": "topLevel", + }, + { + "nodeId": "me-node-1", + "name": "TOPO_NODE_ME4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2", + "nodeType": "TOPO_NODE_ME", + "userLabel": "RMP01", + "xPos": 10, + "yPos": 20, + "parentNode": "66e807c0-94b6-4c50-a02e-fcb56d08bdea", + }, + ], + _Diag(), + ) + + def get_topological_links(self): + return ( + [ + { + "linkId": "415a0bbb-9387-4ff3-8de9-7e36846f4b6a", + "name": "TL{/ME{4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2},EQ={/r=0/sh=1/sl=1},PTP={/p=1_16}_/ME{7508cb6c-59f6-45aa-9e62-4fda61d80553},EQ={/r=0/sh=0/sl=1/ssl=0},PTP={/p=1_20}}", + "connection-status": "Connected", + "owner": "ZTE", + "aEndTpRefList": [ + "ME{4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2},EQ={/r=0/sh=1/sl=1},PTP={/p=1_16}" + ], + "userLabel": "sample-label", + "direction": "BI", + "layerRate": 113, + "zEndTpRefList": [ + "ME{7508cb6c-59f6-45aa-9e62-4fda61d80553},EQ={/r=0/sh=0/sl=1/ssl=0},PTP={/p=1_20}" + ], + }, + { + "linkId": "link-to-drop", + "name": "TL-drop", + "aEndTpRefList": ["ME{aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa},EQ={},PTP={/p=1}"], + "zEndTpRefList": ["ME{bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb},EQ={},PTP={/p=2}"], + "direction": "BI", + "layerRate": 1, + }, + ], + _Diag(), + ) + + class _CNarrow: + def get_topo_nodes(self): + return ( + [ + { + "nodeId": "me-node-1", + "name": "TOPO_NODE_ME4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2", + "nodeType": "TOPO_NODE_ME", + "userLabel": "RMP01-upd", + "xPos": 11, + "yPos": 21, + "parentNode": "topLevel", + } + ], + _Diag(), + ) + + def get_topological_links(self): + return ( + [ + { + "linkId": "415a0bbb-9387-4ff3-8de9-7e36846f4b6a", + "name": "TL-keep", + "connection-status": "Connected", + "aEndTpRefList": [ + "ME{4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2},EQ={/r=0/sh=1/sl=1},PTP={/p=1_16}" + ], + "zEndTpRefList": [ + "ME{7508cb6c-59f6-45aa-9e62-4fda61d80553},EQ={/r=0/sh=0/sl=1/ssl=0},PTP={/p=1_20}" + ], + "direction": "BI", + "layerRate": 113, + } + ], + _Diag(), + ) + + job1 = sync_topology_full(self.db, _CWide(), trigger_mode="manual") + self.assertEqual(job1.status, "done") + self.assertEqual(job1.pulled_count, 4) + self.assertEqual(job1.inserted_count, 4) + + sbn = self.db.get(UmeTopoNode, "66e807c0-94b6-4c50-a02e-fcb56d08bdea") + self.assertIsNotNone(sbn) + self.assertEqual(sbn.node_type, "TOPO_NODE_SBN") + self.assertEqual(sbn.x_pos, 460) + self.assertEqual(sbn.ume_ne_id, "") + + me = self.db.get(UmeTopoNode, "me-node-1") + self.assertIsNotNone(me) + self.assertEqual(me.ume_ne_id, "4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2") + + link = self.db.get(UmeTopoLink, "415a0bbb-9387-4ff3-8de9-7e36846f4b6a") + self.assertIsNotNone(link) + self.assertEqual(link.a_ume_ne_id, "4e598e5d-fe42-4c79-9f62-7d3e5d4eb5b2") + self.assertEqual(link.z_ume_ne_id, "7508cb6c-59f6-45aa-9e62-4fda61d80553") + self.assertEqual(link.a_ptp, "/p=1_16") + self.assertEqual(link.z_ptp, "/p=1_20") + self.assertEqual(link.connection_status, "Connected") + self.assertIsNotNone(self.db.get(UmeTopoLink, "link-to-drop")) + + job2 = sync_topology_full(self.db, _CNarrow(), trigger_mode="manual") + self.assertEqual(job2.status, "done") + self.db.expire_all() + self.assertIsNone(self.db.get(UmeTopoNode, "66e807c0-94b6-4c50-a02e-fcb56d08bdea")) + self.assertIsNotNone(self.db.get(UmeTopoNode, "me-node-1")) + self.assertEqual(self.db.get(UmeTopoNode, "me-node-1").user_label, "RMP01-upd") + self.assertIsNone(self.db.get(UmeTopoLink, "link-to-drop")) + self.assertIsNotNone(self.db.get(UmeTopoLink, "415a0bbb-9387-4ff3-8de9-7e36846f4b6a")) + details = json.loads(job2.details_json or "{}") + self.assertEqual(details.get("deleted_topo_nodes"), 1) + self.assertEqual(details.get("deleted_topo_links"), 1) + + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index d07f4f4..97b4e6f 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -1055,6 +1055,7 @@ const en = { pulling_alarms_current: "Pulling UME current alarms…", alarms_sync_in_progress_skip: "Another current-alarm REST sync in progress; skipped", pulling_inventory: "Pulling UME inventory…", + pulling_topology: "Pulling UME topology (nodes/links)…", ume_ws_disabled_no_base_url: "Disabled or UME_BASE_URL not configured", oclaw_fwd_disabled: "Disabled or NETX_OCLAW_ALARM_WS / token / url not configured", resumed_sync_soon: "Resumed: skipping debounce wait, sync soon", @@ -1092,6 +1093,7 @@ const en = { alarms_current_ws_consumer: "UME alarm WSS", oclaw_alarm_forwarder: "OClaw key-alert WSS", inventory_auto_sync: "Inventory auto sync", + topology_auto_sync: "Topology auto sync", }, }, syncStatus: { diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 75b685d..50bf60c 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -1048,6 +1048,7 @@ const zh = { pulling_alarms_current: "正在拉取 UME 当前告警…", alarms_sync_in_progress_skip: "另一条当前告警 REST 同步进行中,已跳过", pulling_inventory: "正在拉取 UME 网元清单…", + pulling_topology: "正在拉取 UME 拓扑(节点/链路)…", ume_ws_disabled_no_base_url: "未启用或未配置 UME_BASE_URL", oclaw_fwd_disabled: "未启用或未配置 NETX_OCLAW_ALARM_WS / token / url", resumed_sync_soon: "已恢复:将跳过本轮周期等待并尽快同步", @@ -1085,6 +1086,7 @@ const zh = { alarms_current_ws_consumer: "UME 告警 WSS", oclaw_alarm_forwarder: "OClaw 关键告警 WSS", inventory_auto_sync: "Inventory 定时同步", + topology_auto_sync: "拓扑定时同步", }, }, syncStatus: { diff --git a/web/src/utils/runtimeMessages.ts b/web/src/utils/runtimeMessages.ts index 78d9266..f24346a 100644 --- a/web/src/utils/runtimeMessages.ts +++ b/web/src/utils/runtimeMessages.ts @@ -9,6 +9,7 @@ const LEGACY_ZH_RUNTIME_ERROR: Record = { "正在拉取 UME 当前告警…": "ume.tasks.runtimeError.pulling_alarms_current", "另一条当前告警 REST 同步进行中,已跳过": "ume.tasks.runtimeError.alarms_sync_in_progress_skip", "正在拉取 UME 网元清单…": "ume.tasks.runtimeError.pulling_inventory", + "正在拉取 UME 拓扑(节点/链路)…": "ume.tasks.runtimeError.pulling_topology", "未启用或未配置 UME_BASE_URL": "ume.tasks.runtimeError.ume_ws_disabled_no_base_url", "未启用或未配置 NETX_OCLAW_ALARM_WS / token / url": "ume.tasks.runtimeError.oclaw_fwd_disabled", "已恢复:将跳过本轮周期等待并尽快同步": "ume.tasks.runtimeError.resumed_sync_soon",