From 340fc8ee7dc2b14f4692de1cf31f3a9229cabc5a Mon Sep 17 00:00:00 2001 From: oliver Date: Mon, 11 May 2026 18:31:04 +0800 Subject: [PATCH] Reconcile UME inventory and current alarms with full snapshot After a complete marker pull (is_end_of_reply, no iterator/dup abort), delete local NE/holders and current alarms not present in this batch. Replace 48h-only current-alarm expiry. Flush before bulk delete to avoid ORM stale updates. Job/batch JSON records deleted counts and reconcile flags. Co-authored-by: Cursor --- netx_api/ume_sync_service.py | 80 +++++++++++++++++++++---- tests/test_ume_sync.py | 112 +++++++++++++++++++++++++++++++++++ 2 files changed, 180 insertions(+), 12 deletions(-) diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 26e7f85..5180ee9 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -1,7 +1,7 @@ from __future__ import annotations import json -from datetime import datetime, timedelta, timezone +from datetime import datetime, timezone from typing import Any from sqlalchemy.orm import Session @@ -80,6 +80,22 @@ def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: ) +def _snapshot_reconcile_ok(meta: dict[str, Any]) -> bool: + """True when paging finished normally (full snapshot); avoid deleting local rows on partial pulls.""" + if not bool(meta.get("is_end_of_reply")): + return False + if bool(meta.get("graceful_end_by_iterator_error")): + return False + warnings = meta.get("warnings") or [] + if not isinstance(warnings, list): + return False + if "duplicate_page_detected" in [str(w) for w in warnings]: + return False + if str(meta.get("paging_note") or "").strip() == "duplicate_page_detected": + return False + return True + + def _collect_marker_pages( fetch_page: Any, *, @@ -162,7 +178,7 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " max_pages = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000) max_pages = max(1, min(max_pages, 20000)) - pages, _ = _collect_marker_pages( + pages, inv_meta = _collect_marker_pages( lambda marker: client.get_network_elements(limit=page_size, marker=marker), max_pages=max_pages, iterator_500_as_end=False, @@ -170,10 +186,12 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " ne_rows = [row for page in pages for row in page] now = _utc_now_naive() pulled = len(ne_rows) + seen_ne_ids: set[str] = set() for row in ne_rows: ne_id = _s(_pick(row, "ne-id", "ne_id", "id")) if not ne_id: continue + seen_ne_ids.add(ne_id) existing = db.get(UmeInventoryNE, ne_id) if existing is None: existing = UmeInventoryNE( @@ -245,6 +263,37 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " existing_holder.last_seen_at = now existing_holder.raw_json = json.dumps(h, ensure_ascii=False, default=str) + db.flush() + deleted_holders = deleted_ne = 0 + if _snapshot_reconcile_ok(inv_meta): + if seen_ne_ids: + deleted_holders = int( + db.query(UmeInventoryEquipmentHolder) + .filter(~UmeInventoryEquipmentHolder.ne_id.in_(list(seen_ne_ids))) + .delete(synchronize_session=False) + ) + deleted_ne = int( + db.query(UmeInventoryNE) + .filter(~UmeInventoryNE.ne_id.in_(list(seen_ne_ids))) + .delete(synchronize_session=False) + ) + else: + deleted_holders = int(db.query(UmeInventoryEquipmentHolder).delete(synchronize_session=False)) + deleted_ne = int(db.query(UmeInventoryNE).delete(synchronize_session=False)) + + job.details_json = json.dumps( + { + "inventory_reconcile": _snapshot_reconcile_ok(inv_meta), + "deleted_inventory_holders": deleted_holders, + "deleted_inventory_ne": deleted_ne, + "paging": { + "is_end_of_reply": bool(inv_meta.get("is_end_of_reply")), + "warnings": list(inv_meta.get("warnings") or []), + }, + }, + ensure_ascii=False, + ) + job.status = "done" except Exception as exc: job.status = "failed" @@ -276,8 +325,8 @@ def _sync_alarms_common( ) db.add(batch) db.flush() - now = _utc_now_naive() pulled = inserted = updated = 0 + deleted_stale_current = 0 paging_mode = "marker" paging_note = "" page_no = 0 @@ -285,6 +334,7 @@ def _sync_alarms_common( is_end_of_reply = False graceful_end_by_iterator_error = False warnings: list[str] = [] + meta: dict[str, Any] = {} try: limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = max(1, limit_max) @@ -292,13 +342,14 @@ def _sync_alarms_common( page_size = max(1, min(page_size, limit_max)) max_pages = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000) max_pages = max(1, min(max_pages, 20000)) - def upsert_alarm(alarm: dict[str, Any]) -> None: + + def upsert_alarm(alarm: dict[str, Any], *, touch_ts: datetime) -> None: nonlocal inserted, updated key = _alarm_key(alarm) if is_uncleared: existing = db.get(UmeAlarmHistory, key) if existing is None: - existing = UmeAlarmHistory(alarm_key=key, first_seen_at=now) + existing = UmeAlarmHistory(alarm_key=key, first_seen_at=touch_ts) db.add(existing) inserted += 1 else: @@ -306,7 +357,7 @@ def _sync_alarms_common( else: existing = db.get(UmeAlarmCurrent, key) if existing is None: - existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=now) + existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=touch_ts) db.add(existing) inserted += 1 else: @@ -321,7 +372,7 @@ def _sync_alarms_common( existing.root_cause_alarm_indication = _s( _pick(alarm, "rootCauseAlarmIndication", "root-cause-alarm-indication") ) - existing.last_seen_at = now + existing.last_seen_at = touch_ts existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str) iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True)) @@ -330,10 +381,12 @@ def _sync_alarms_common( max_pages=max_pages, iterator_500_as_end=iterator_500_as_end, ) + sync_batch_ts = _utc_now_naive() for rows in pages: pulled += len(rows) for alarm in rows: - upsert_alarm(alarm) + upsert_alarm(alarm, touch_ts=sync_batch_ts) + db.flush() page_no = int(meta.get("page_count") or 0) next_marker = str(meta.get("last_marker") or "") is_end_of_reply = bool(meta.get("is_end_of_reply")) @@ -341,11 +394,10 @@ def _sync_alarms_common( paging_note = str(meta.get("paging_note") or "") warnings = [str(x) for x in (meta.get("warnings") or []) if str(x)] - if not is_uncleared: - expiry = now - timedelta(hours=48) - ( + if not is_uncleared and _snapshot_reconcile_ok(meta): + deleted_stale_current = int( db.query(UmeAlarmCurrent) - .filter(UmeAlarmCurrent.last_seen_at < expiry) + .filter(UmeAlarmCurrent.last_seen_at < sync_batch_ts) .delete(synchronize_session=False) ) @@ -365,6 +417,8 @@ def _sync_alarms_common( "is_end_of_reply": is_end_of_reply, "graceful_end_by_iterator_error": graceful_end_by_iterator_error, "warnings": warnings, + "deleted_stale_current_alarms": int(deleted_stale_current), + "current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta), }, ensure_ascii=False, ) @@ -394,6 +448,8 @@ def _sync_alarms_common( "is_end_of_reply": is_end_of_reply, "graceful_end_by_iterator_error": graceful_end_by_iterator_error, "warnings": warnings, + "deleted_stale_current_alarms": int(deleted_stale_current), + "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 f1df96b..1a74453 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -207,6 +207,118 @@ class UmeSyncServiceTests(unittest.TestCase): self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-1")) self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-2")) + def test_sync_inventory_reconcile_removes_ne_not_in_snapshot(self): + class _CWide: + def get_network_elements(self, *, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = True + + return ( + [ + {"ne-id": "NE-1", "name": "ne1"}, + {"ne-id": "NE-2", "name": "ne2"}, + ], + _D(), + ) + + class _CNarrow: + def get_network_elements(self, *, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = True + + return ([{"ne-id": "NE-1", "name": "ne1"}], _D()) + + sync_inventory_full(self.db, _CWide(), trigger_mode="manual") + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-2")) + + sync_inventory_full(self.db, _CNarrow(), trigger_mode="manual") + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-1")) + self.assertIsNone(self.db.get(UmeInventoryNE, "NE-2")) + + def test_sync_inventory_partial_pull_does_not_delete(self): + class _CPartial: + def get_network_elements(self, *, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = False + + return ([{"ne-id": "NE-9", "name": "ne9"}], _D()) + + self.db.add(UmeInventoryNE(ne_id="NE-OLD", ne_name="old")) + self.db.commit() + + sync_inventory_full(self.db, _CPartial(), trigger_mode="manual") + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-9")) + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-OLD")) + + def test_sync_current_alarms_reconcile_drops_missing_keys(self): + class _COne: + def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = True + + return ( + [ + { + "alarmKey": "AK-KEEP", + "ne-id": "NE-1", + "perceivedSeverity": "major", + "isCleared": "false", + }, + ], + _D(), + ) + + self.db.add( + UmeAlarmCurrent( + alarm_key="AK-GONE", + ne_id="NE-X", + perceived_severity="minor", + is_cleared="false", + ) + ) + self.db.commit() + + sync_alarms_current(self.db, _COne(), trigger_mode="manual") + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-KEEP")) + self.assertIsNone(self.db.get(UmeAlarmCurrent, "AK-GONE")) + + def test_sync_current_alarms_partial_pull_keeps_stale_rows(self): + class _CPartial: + def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = False + + return ( + [ + { + "alarmKey": "AK-NEW", + "ne-id": "NE-1", + "perceivedSeverity": "major", + "isCleared": "false", + }, + ], + _D(), + ) + + self.db.add( + UmeAlarmCurrent( + alarm_key="AK-STALE", + ne_id="NE-X", + perceived_severity="minor", + is_cleared="false", + ) + ) + self.db.commit() + + sync_alarms_current(self.db, _CPartial(), trigger_mode="manual") + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-NEW")) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-STALE")) + def test_sync_current_alarms_upsert(self): class _C: def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None):