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 <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-11 18:31:04 +08:00
parent c5f82f512d
commit 340fc8ee7d
2 changed files with 180 additions and 12 deletions

View file

@ -1,7 +1,7 @@
from __future__ import annotations from __future__ import annotations
import json import json
from datetime import datetime, timedelta, timezone from datetime import datetime, timezone
from typing import Any from typing import Any
from sqlalchemy.orm import Session 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( def _collect_marker_pages(
fetch_page: Any, 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 = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000)
max_pages = max(1, min(max_pages, 20000)) 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), lambda marker: client.get_network_elements(limit=page_size, marker=marker),
max_pages=max_pages, max_pages=max_pages,
iterator_500_as_end=False, 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] ne_rows = [row for page in pages for row in page]
now = _utc_now_naive() now = _utc_now_naive()
pulled = len(ne_rows) pulled = len(ne_rows)
seen_ne_ids: set[str] = set()
for row in ne_rows: for row in ne_rows:
ne_id = _s(_pick(row, "ne-id", "ne_id", "id")) ne_id = _s(_pick(row, "ne-id", "ne_id", "id"))
if not ne_id: if not ne_id:
continue continue
seen_ne_ids.add(ne_id)
existing = db.get(UmeInventoryNE, ne_id) existing = db.get(UmeInventoryNE, ne_id)
if existing is None: if existing is None:
existing = UmeInventoryNE( 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.last_seen_at = now
existing_holder.raw_json = json.dumps(h, ensure_ascii=False, default=str) 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" job.status = "done"
except Exception as exc: except Exception as exc:
job.status = "failed" job.status = "failed"
@ -276,8 +325,8 @@ def _sync_alarms_common(
) )
db.add(batch) db.add(batch)
db.flush() db.flush()
now = _utc_now_naive()
pulled = inserted = updated = 0 pulled = inserted = updated = 0
deleted_stale_current = 0
paging_mode = "marker" paging_mode = "marker"
paging_note = "" paging_note = ""
page_no = 0 page_no = 0
@ -285,6 +334,7 @@ def _sync_alarms_common(
is_end_of_reply = False is_end_of_reply = False
graceful_end_by_iterator_error = False graceful_end_by_iterator_error = False
warnings: list[str] = [] warnings: list[str] = []
meta: dict[str, Any] = {}
try: try:
limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000)
limit_max = max(1, limit_max) limit_max = max(1, limit_max)
@ -292,13 +342,14 @@ def _sync_alarms_common(
page_size = max(1, min(page_size, limit_max)) 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 = int(getattr(settings, "ume_marker_max_pages", getattr(settings, "ume_max_pages", 2000)) or 2000)
max_pages = max(1, min(max_pages, 20000)) 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 nonlocal inserted, updated
key = _alarm_key(alarm) key = _alarm_key(alarm)
if is_uncleared: if is_uncleared:
existing = db.get(UmeAlarmHistory, key) existing = db.get(UmeAlarmHistory, key)
if existing is None: 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) db.add(existing)
inserted += 1 inserted += 1
else: else:
@ -306,7 +357,7 @@ def _sync_alarms_common(
else: else:
existing = db.get(UmeAlarmCurrent, key) existing = db.get(UmeAlarmCurrent, key)
if existing is None: 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) db.add(existing)
inserted += 1 inserted += 1
else: else:
@ -321,7 +372,7 @@ def _sync_alarms_common(
existing.root_cause_alarm_indication = _s( existing.root_cause_alarm_indication = _s(
_pick(alarm, "rootCauseAlarmIndication", "root-cause-alarm-indication") _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) existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str)
iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True)) 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, max_pages=max_pages,
iterator_500_as_end=iterator_500_as_end, iterator_500_as_end=iterator_500_as_end,
) )
sync_batch_ts = _utc_now_naive()
for rows in pages: for rows in pages:
pulled += len(rows) pulled += len(rows)
for alarm in 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) page_no = int(meta.get("page_count") or 0)
next_marker = str(meta.get("last_marker") or "") next_marker = str(meta.get("last_marker") or "")
is_end_of_reply = bool(meta.get("is_end_of_reply")) 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 "") paging_note = str(meta.get("paging_note") or "")
warnings = [str(x) for x in (meta.get("warnings") or []) if str(x)] warnings = [str(x) for x in (meta.get("warnings") or []) if str(x)]
if not is_uncleared: if not is_uncleared and _snapshot_reconcile_ok(meta):
expiry = now - timedelta(hours=48) deleted_stale_current = int(
(
db.query(UmeAlarmCurrent) db.query(UmeAlarmCurrent)
.filter(UmeAlarmCurrent.last_seen_at < expiry) .filter(UmeAlarmCurrent.last_seen_at < sync_batch_ts)
.delete(synchronize_session=False) .delete(synchronize_session=False)
) )
@ -365,6 +417,8 @@ def _sync_alarms_common(
"is_end_of_reply": is_end_of_reply, "is_end_of_reply": is_end_of_reply,
"graceful_end_by_iterator_error": graceful_end_by_iterator_error, "graceful_end_by_iterator_error": graceful_end_by_iterator_error,
"warnings": warnings, "warnings": warnings,
"deleted_stale_current_alarms": int(deleted_stale_current),
"current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta),
}, },
ensure_ascii=False, ensure_ascii=False,
) )
@ -394,6 +448,8 @@ def _sync_alarms_common(
"is_end_of_reply": is_end_of_reply, "is_end_of_reply": is_end_of_reply,
"graceful_end_by_iterator_error": graceful_end_by_iterator_error, "graceful_end_by_iterator_error": graceful_end_by_iterator_error,
"warnings": warnings, "warnings": warnings,
"deleted_stale_current_alarms": int(deleted_stale_current),
"current_snapshot_reconcile": (not is_uncleared) and _snapshot_reconcile_ok(meta),
}, },
ensure_ascii=False, ensure_ascii=False,
) )

View file

@ -207,6 +207,118 @@ class UmeSyncServiceTests(unittest.TestCase):
self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-1")) self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-1"))
self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-2")) 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): def test_sync_current_alarms_upsert(self):
class _C: class _C:
def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None):