diff --git a/.env.example b/.env.example index 37764a8..0a7b75e 100644 --- a/.env.example +++ b/.env.example @@ -32,7 +32,8 @@ 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_TOPOLOGY_TIMEOUT_S=600 +# NETX_UME_SYNC_TOPOLOGY_STALE_RUNNING_SEC=1800 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 07a81a9..d814a50 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -58,7 +58,8 @@ class Settings(BaseSettings): 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_topology_timeout_s: float = 600.0 + ume_sync_topology_stale_running_sec: int = 1800 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" diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index 4116ffd..799fc4c 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -241,7 +241,12 @@ class UMEClient: # 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) - to = self.timeout_s if timeout_s is None else max(3.0, float(timeout_s)) + if timeout_s is None: + to: float | httpx.Timeout = self.timeout_s + else: + # Large topology dumps need a long read window; connect stays short. + read_s = max(3.0, float(timeout_s)) + to = httpx.Timeout(connect=min(30.0, read_s), read=read_s, write=min(60.0, read_s), pool=30.0) return httpx.Client(transport=transport, timeout=to) @@ -615,11 +620,38 @@ class UMEClient: 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) + return max(30.0, base) + return max(self.timeout_s, 600.0) + + @staticmethod + def _extract_restconf_list(payload: dict[str, Any], *, containers: list[str], items: list[str]) -> list[dict[str, Any]]: + """Fast path for TopoNodes/TopologicalLinks without O(n) json.dumps dedupe walks.""" + container_keys = {str(k).lower() for k in containers} + item_keys = {str(k).lower() for k in items} + for ck, cv in payload.items(): + if str(ck).lower() not in container_keys: + continue + if isinstance(cv, list): + return [x for x in cv if isinstance(x, dict)] + if isinstance(cv, dict): + for ik, iv in cv.items(): + if str(ik).lower() in item_keys and isinstance(iv, list): + return [x for x in iv if isinstance(x, dict)] + # Some payloads nest once more under the same item key. + for iv in cv.values(): + if isinstance(iv, list) and iv and isinstance(iv[0], dict): + return [x for x in iv if isinstance(x, dict)] + return [] 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_restconf_list( + data, + containers=["TopoNodes", "topo-nodes", "topoNodes"], + items=["TopoNode", "topo-node", "topoNode"], + ) + if rows: + return rows, diag rows = self._extract_named_list(data, ["TopoNodes", "TopoNode", "topo-nodes", "topo-node"]) if rows: return rows, diag @@ -633,6 +665,13 @@ class UMEClient: data, diag = self.request_json( "GET", self.topological_links_path, timeout_s=self._topology_timeout_s() ) + rows = self._extract_restconf_list( + data, + containers=["TopologicalLinks", "topological-links", "topologicalLinks"], + items=["TopologicalLink", "topological-link", "topologicalLink"], + ) + if rows: + return rows, diag rows = self._extract_named_list( data, ["TopologicalLinks", "TopologicalLink", "topological-links", "topological-link"] ) diff --git a/netx_api/ume_runtime.py b/netx_api/ume_runtime.py index 3625fa0..e6bed21 100644 --- a/netx_api/ume_runtime.py +++ b/netx_api/ume_runtime.py @@ -25,6 +25,7 @@ from .ume_alarm_ws import ( start_ume_alarm_ws_consumer, ) from .ume_sync_service import sync_alarms_current, sync_inventory_full, sync_topology_full +from .ume_sync_topology import fail_stale_topology_running_jobs from .runtime_task_messages import ( RT_ALARMS_SYNC_IN_PROGRESS_SKIP, RT_KEEPALIVE_FAILED, @@ -323,6 +324,7 @@ def start_api_sideband_threads() -> None: ) db = SessionLocal() try: + fail_stale_topology_running_jobs(db) client = ume_support._ume_client() sync_topology_full(db, client, trigger_mode="schedule") _schedule_log.info("topology_auto_sync: sync finished ok") @@ -334,6 +336,17 @@ def start_api_sideband_threads() -> None: ) finally: db.close() + except RuntimeError as exc: + if str(exc).startswith("topology_sync_busy"): + _schedule_log.info("topology_auto_sync: skipped busy (%s)", exc) + ume_support._refresh_runtime_task_idle( + "topology_auto_sync", + "topology", + last_error="rt:topology_sync_busy", + ) + time.sleep(30) + else: + raise except Exception as exc: _schedule_log.exception("topology_auto_sync: sync failed: %s", exc) ume_support._set_runtime_task( diff --git a/netx_api/ume_sync_router.py b/netx_api/ume_sync_router.py index b30be78..a5d5120 100644 --- a/netx_api/ume_sync_router.py +++ b/netx_api/ume_sync_router.py @@ -66,17 +66,33 @@ def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db } ) 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 ""), - } - ) + try: + 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 ""), + } + ) + except RuntimeError as exc: + msg = str(exc) + if msg.startswith("topology_sync_busy"): + out["jobs"].append( + { + "domain": "topology", + "status": "skipped", + "pulled_count": 0, + "inserted_count": 0, + "updated_count": 0, + "error_message": msg[:240], + } + ) + else: + raise 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": diff --git a/netx_api/ume_sync_topology.py b/netx_api/ume_sync_topology.py index 7c25453..9fa57db 100644 --- a/netx_api/ume_sync_topology.py +++ b/netx_api/ume_sync_topology.py @@ -5,10 +5,15 @@ import hashlib import json import logging import re +import threading +from datetime import timedelta from typing import Any -from sqlalchemy.orm import Session +from sqlalchemy.exc import IntegrityError, OperationalError, InvalidRequestError +from sqlalchemy.orm import Session, sessionmaker +from .config import settings +from .db import SessionLocal from .models import UmeSyncJob, UmeTopoLink, UmeTopoNode from .ume_client import UMEClient from .ume_raw import dumps_ume_raw @@ -16,6 +21,7 @@ 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") +_TOPOLOGY_SYNC_LOCK = threading.Lock() _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}" @@ -35,7 +41,6 @@ def extract_me_uuid(text: str) -> str: 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") @@ -103,161 +108,437 @@ def _ume_ne_id_for_topo_node(*, node_type: str, name: str) -> str: return extract_me_uuid(name) +def _stale_running_sec() -> int: + return max(300, int(getattr(settings, "ume_sync_topology_stale_running_sec", 1800) or 1800)) + + +def fail_stale_topology_running_jobs(db: Session | None = None) -> int: + """Close topology jobs stuck in running (hung HTTP / dead session finalize).""" + own = db is None + sess = db if db is not None else SessionLocal() + try: + cutoff = _utc_now_naive() - timedelta(seconds=_stale_running_sec()) + rows = ( + sess.query(UmeSyncJob) + .filter( + UmeSyncJob.domain == "topology", + UmeSyncJob.status == "running", + UmeSyncJob.ended_at.is_(None), + UmeSyncJob.started_at < cutoff, + ) + .all() + ) + if not rows: + return 0 + now = _utc_now_naive() + for row in rows: + row.status = "failed" + row.ended_at = now + msg = str(row.error_message or "").strip() + suffix = "stale_running_topology_reaped" + row.error_message = (msg + ("; " if msg else "") + suffix)[:1024] + sess.commit() + _sync_log.warning("topology sync: reaped %s stale running jobs", len(rows)) + return len(rows) + except Exception: + _sync_log.exception("topology sync: stale reap failed") + try: + sess.rollback() + except Exception: + pass + return 0 + finally: + if own: + sess.close() + + +def _active_topology_running(db: Session) -> UmeSyncJob | None: + cutoff = _utc_now_naive() - timedelta(seconds=_stale_running_sec()) + return ( + db.query(UmeSyncJob) + .filter( + UmeSyncJob.domain == "topology", + UmeSyncJob.status == "running", + UmeSyncJob.ended_at.is_(None), + UmeSyncJob.started_at >= cutoff, + ) + .order_by(UmeSyncJob.id.desc()) + .first() + ) + + +def _session_factory_for(db: Session): + """Open a sibling session on the same bind (tests use sqlite; prod uses pool).""" + bind = db.get_bind() + return sessionmaker(bind=bind, autoflush=False, autocommit=False, expire_on_commit=False) + + +def _finalize_topology_job( + job_id: int, + *, + status: str, + pulled: int, + inserted: int, + updated: int, + error_message: str = "", + details_json: str = "{}", + db: Session | None = None, +) -> UmeSyncJob | None: + """Persist job end state on a fresh session so a dead pull-session cannot leave running forever.""" + if db is not None: + sess = _session_factory_for(db)() + else: + sess = SessionLocal() + try: + row = sess.get(UmeSyncJob, int(job_id)) + if row is None: + return None + row.status = status + row.pulled_count = int(pulled) + row.inserted_count = int(inserted) + row.updated_count = int(updated) + row.error_message = str(error_message or "")[:1024] + row.details_json = str(details_json or "{}") + row.ended_at = _utc_now_naive() + sess.commit() + sess.refresh(row) + return row + except Exception: + _sync_log.exception("topology sync: finalize job %s failed", job_id) + try: + sess.rollback() + except Exception: + pass + return None + finally: + sess.close() + + +def _safe_rollback(sess: Session) -> None: + try: + sess.rollback() + except Exception: + pass + + +def _dedupe_rows_by_id( + rows: list[Any], id_fn +) -> tuple[list[tuple[str, dict[str, Any]]], int]: + """Keep last row per id; drop non-dicts / missing ids.""" + by_id: dict[str, dict[str, Any]] = {} + skipped = 0 + for row in rows: + if not isinstance(row, dict): + skipped += 1 + continue + rid = str(id_fn(row) or "").strip() + if not rid: + skipped += 1 + continue + by_id[rid] = row + return list(by_id.items()), skipped + + +def _load_node(work: Session, node_id: str) -> UmeTopoNode | None: + row = work.get(UmeTopoNode, node_id) + if row is not None: + return row + return work.query(UmeTopoNode).filter(UmeTopoNode.node_id == node_id).one_or_none() + + +def _load_link(work: Session, link_id: str) -> UmeTopoLink | None: + row = work.get(UmeTopoLink, link_id) + if row is not None: + return row + return work.query(UmeTopoLink).filter(UmeTopoLink.link_id == link_id).one_or_none() + + +def _upsert_topo_node(work: Session, node_id: str, row: dict[str, Any], *, now) -> str: + """Return 'inserted' | 'updated'.""" + name = _s(_pick(row, "name")) + node_type = _s(_pick(row, "nodeType", "node-type", "node_type")) + + def _apply(existing: UmeTopoNode) -> None: + 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) + + existing = _load_node(work, node_id) + if existing is not None: + _apply(existing) + return "updated" + + existing = UmeTopoNode(node_id=node_id, first_seen_at=now) + _apply(existing) + try: + with work.begin_nested(): + work.add(existing) + work.flush() + return "inserted" + except IntegrityError: + try: + work.expunge(existing) + except Exception: + pass + existing = _load_node(work, node_id) + if existing is None: + raise + _apply(existing) + return "updated" + + +def _upsert_topo_link(work: Session, link_id: str, row: dict[str, Any], *, now) -> str: + """Return 'inserted' | 'updated'.""" + 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")) + + def _apply(existing: UmeTopoLink) -> None: + 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) + + existing = _load_link(work, link_id) + if existing is not None: + _apply(existing) + return "updated" + + existing = UmeTopoLink(link_id=link_id, first_seen_at=now) + _apply(existing) + try: + with work.begin_nested(): + work.add(existing) + work.flush() + return "inserted" + except IntegrityError: + try: + work.expunge(existing) + except Exception: + pass + existing = _load_link(work, link_id) + if existing is None: + raise + _apply(existing) + return "updated" + + +def _delete_missing_ids(work: Session, model, pk_col, seen_ids: set[str]) -> int: + if not seen_ids: + return int(work.query(model).delete(synchronize_session=False)) + all_pks = [str(x[0]) for x in work.query(pk_col).all() if str(x[0] or "").strip()] + missing = [pk for pk in all_pks if pk not in seen_ids] + deleted = 0 + chunk = 5000 + for i in range(0, len(missing), chunk): + part = missing[i : i + chunk] + deleted += int(work.query(model).filter(pk_col.in_(part)).delete(synchronize_session=False)) + return deleted + + 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 + if not _TOPOLOGY_SYNC_LOCK.acquire(blocking=False): + raise RuntimeError("topology_sync_busy") 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) + fail_stale_topology_running_jobs(db) + active = _active_topology_running(db) + if active is not None: + raise RuntimeError(f"topology_sync_busy:job_id={active.id}") + job = _build_sync_job("topology", trigger_mode) + db.add(job) 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 + job_id = int(job.id) + _sync_log.info( + "topology sync job %s committed as running (trigger=%s)", + 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 + details = "{}" + err = "" + status = "done" + + try: + try: + db.expire_all() + except Exception: + pass + + work = _session_factory_for(db)() + try: + now = _utc_now_naive() + _sync_log.info("topology sync job %s: pulling TopoNodes…", job_id) + node_rows, node_diag = client.get_topo_nodes() + nodes_pulled = len(node_rows) + pairs, nodes_skip = _dedupe_rows_by_id(node_rows, _node_id_from_row) + _sync_log.info( + "topology sync job %s: TopoNodes rows=%s unique=%s latency_ms=%s", + job_id, + nodes_pulled, + len(pairs), + int(getattr(node_diag, "latency_ms", 0) or 0), + ) + seen_nodes: set[str] = set() + for i, (node_id, row) in enumerate(pairs): + seen_nodes.add(node_id) + try: + action = _upsert_topo_node(work, node_id, row, now=now) + except (IntegrityError, OperationalError, InvalidRequestError): + _safe_rollback(work) + action = _upsert_topo_node(work, node_id, row, now=now) + if action == "inserted": + nodes_ins += 1 + else: + nodes_upd += 1 + if (i + 1) % 2000 == 0: + work.commit() + _sync_log.info( + "topology sync job %s: nodes progress %s/%s", + job_id, + i + 1, + len(pairs), + ) + + work.flush() + nodes_del = _delete_missing_ids(work, UmeTopoNode, UmeTopoNode.node_id, seen_nodes) + work.commit() + + _sync_log.info("topology sync job %s: pulling TopologicalLinks…", job_id) + link_rows, link_diag = client.get_topological_links() + links_pulled = len(link_rows) + link_pairs, links_skip = _dedupe_rows_by_id(link_rows, _link_id_from_row) + _sync_log.info( + "topology sync job %s: TopologicalLinks rows=%s unique=%s latency_ms=%s", + job_id, + links_pulled, + len(link_pairs), + int(getattr(link_diag, "latency_ms", 0) or 0), + ) + seen_links: set[str] = set() + for i, (link_id, row) in enumerate(link_pairs): + seen_links.add(link_id) + try: + action = _upsert_topo_link(work, link_id, row, now=now) + except (IntegrityError, OperationalError, InvalidRequestError): + _safe_rollback(work) + action = _upsert_topo_link(work, link_id, row, now=now) + if action == "inserted": + links_ins += 1 + else: + links_upd += 1 + if (i + 1) % 2000 == 0: + work.commit() + _sync_log.info( + "topology sync job %s: links progress %s/%s", + job_id, + i + 1, + len(link_pairs), + ) + + work.flush() + links_del = _delete_missing_ids(work, UmeTopoLink, UmeTopoLink.link_id, seen_links) + work.commit() + + pulled = nodes_pulled + links_pulled + inserted = nodes_ins + links_ins + updated = nodes_upd + links_upd + details = json.dumps( + { + "nodes": { + "pulled": nodes_pulled, + "unique": len(seen_nodes), + "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, + "unique": len(seen_links), + "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, + ) + _sync_log.info( + "topology sync done job=%s nodes=%s/%s/%s links=%s/%s/%s deleted_n=%s deleted_l=%s", + job_id, + nodes_pulled, + nodes_ins, + nodes_upd, + links_pulled, + links_ins, + links_upd, + nodes_del, + links_del, + ) + finally: + try: + work.close() + except Exception: + pass + except Exception as exc: + status = "failed" + err = str(exc)[:1024] + _sync_log.exception("topology sync failed job=%s: %s", job_id, exc) + + finalized = _finalize_topology_job( + job_id, + status=status, + pulled=pulled, + inserted=inserted, + updated=updated, + error_message=err, + details_json=details, + db=db, + ) + if finalized is not None: + return finalized + orphan = UmeSyncJob( + domain="topology", + status=status, + trigger_mode=trigger_mode, + pulled_count=pulled, + inserted_count=inserted, + updated_count=updated, + error_message=(err or "finalize_failed")[:1024], + details_json=details, + ended_at=_utc_now_naive(), + ) + orphan.id = job_id + return orphan + finally: + _TOPOLOGY_SYNC_LOCK.release() diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index be40c7f..ab9369d 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -1243,6 +1243,49 @@ class UmeTopologySyncTests(unittest.TestCase): self.assertEqual(details.get("deleted_topo_nodes"), 1) self.assertEqual(details.get("deleted_topo_links"), 1) + def test_sync_topology_duplicate_ids_in_payload(self): + class _Diag: + latency_ms = 1 + + class _C: + def get_topo_nodes(self): + row = { + "nodeId": "dup-node", + "name": "TOPO_NODE_MEaaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa", + "nodeType": "TOPO_NODE_ME", + "userLabel": "v1", + "xPos": 1, + "yPos": 2, + } + row2 = dict(row) + row2["userLabel"] = "v2" + row2["xPos"] = 9 + return ([row, row2], _Diag()) + + def get_topological_links(self): + link = { + "linkId": "dup-link", + "name": "TL-1", + "aEndTpRefList": ["ME{aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa},PTP={/p=1}"], + "zEndTpRefList": ["ME{bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb},PTP={/p=2}"], + "direction": "BI", + "layerRate": 1, + } + link2 = dict(link) + link2["layerRate"] = 2 + return ([link, link2], _Diag()) + + job = sync_topology_full(self.db, _C(), trigger_mode="manual") + self.assertEqual(job.status, "done", job.error_message) + self.db.expire_all() + node = self.db.get(UmeTopoNode, "dup-node") + self.assertIsNotNone(node) + self.assertEqual(node.user_label, "v2") + self.assertEqual(node.x_pos, 9) + link = self.db.get(UmeTopoLink, "dup-link") + self.assertIsNotNone(link) + self.assertEqual(link.layer_rate, 2) + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 97b4e6f..51f4379 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -1056,6 +1056,7 @@ const en = { 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)…", + topology_sync_busy: "Topology sync already running; skipped", 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", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 50bf60c..c748304 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -1049,6 +1049,7 @@ const zh = { alarms_sync_in_progress_skip: "另一条当前告警 REST 同步进行中,已跳过", pulling_inventory: "正在拉取 UME 网元清单…", pulling_topology: "正在拉取 UME 拓扑(节点/链路)…", + topology_sync_busy: "拓扑同步进行中,已跳过", ume_ws_disabled_no_base_url: "未启用或未配置 UME_BASE_URL", oclaw_fwd_disabled: "未启用或未配置 NETX_OCLAW_ALARM_WS / token / url", resumed_sync_soon: "已恢复:将跳过本轮周期等待并尽快同步",