From 037873423d724bb6acea71c816f36f97fbae8ece Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 8 May 2026 21:44:21 +0800 Subject: [PATCH] =?UTF-8?q?feat(UME):=20=E7=BD=91=E5=85=83=E4=B8=8E?= =?UTF-8?q?=E5=91=8A=E8=AD=A6=E7=BB=9F=E4=B8=80marker=E5=88=86=E9=A1=B5?= =?UTF-8?q?=E8=83=BD=E5=8A=9B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将 marker 分页逻辑抽象为通用流程并同时应用于网元和告警同步,默认按 limit+marker 拉取全量数据,统一处理 is-end-of-reply、缺失marker和重复页保护。 Co-authored-by: Cursor --- netx_api/ume_client.py | 19 ++++- netx_api/ume_sync_service.py | 148 +++++++++++++++++++++++------------ tests/test_ume_sync.py | 64 ++++++++++++++- 3 files changed, 178 insertions(+), 53 deletions(-) diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index f5f1cb7..b5210e0 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -442,8 +442,23 @@ class UMEClient: walk(payload) return found - def get_network_elements(self) -> tuple[list[dict[str, Any]], RequestDiagnostics]: - data, diag = self.request_json("GET", self.ne_path) + def get_network_elements( + self, + *, + limit: int | None = None, + marker: str | None = None, + ) -> tuple[list[dict[str, Any]], RequestDiagnostics]: + limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) + limit_max = max(1, limit_max) + page_size = int(limit or settings.ume_page_size or 1000) + page_size = max(1, min(page_size, limit_max)) + params: dict[str, Any] = { + "limit": page_size, + } + marker_value = str(marker or "").strip() + if marker_value: + params["marker"] = marker_value + data, diag = self.request_json("GET", self.ne_path, params=params) rows = self._extract_named_list(data, ["network-elements", "network-element", "ne", "network_elements"]) if rows: return rows, diag diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 07ae51e..5bfe8ce 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -80,13 +80,94 @@ def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: ) +def _collect_marker_pages( + fetch_page: Any, + *, + max_pages: int, + iterator_500_as_end: bool = False, +) -> tuple[list[list[dict[str, Any]]], dict[str, Any]]: + page_no = 0 + next_marker = "" + is_end_of_reply = False + graceful_end_by_iterator_error = False + paging_note = "" + warnings: list[str] = [] + last_page_signature = "" + pages: list[list[dict[str, Any]]] = [] + + while True: + page_no += 1 + if page_no > max_pages: + raise RuntimeError(f"ume_alarms_pagination_exceeded:max_pages={max_pages}") + + try: + rows, diag = fetch_page(next_marker or None) + except Exception as exc: + msg = str(exc or "") + low = msg.lower() + if iterator_500_as_end and pages and "ume_request_failed:500" in low and "iterator" in low and "null" in low: + graceful_end_by_iterator_error = True + paging_note = msg[:240] + break + raise + + rows = [x for x in rows if isinstance(x, dict)] + pages.append(rows) + + # Protection against repeated pages causing infinite loops. + cur_sig = "|".join(sorted(_alarm_key(x) for x in rows)) + if cur_sig and cur_sig == last_page_signature: + warnings.append("duplicate_page_detected") + paging_note = "duplicate_page_detected" + break + last_page_signature = cur_sig + + has_is_end = diag.is_end_of_reply is not None + is_end_of_reply = bool(diag.is_end_of_reply) if has_is_end else False + next_marker = str(diag.marker or "").strip() + + if has_is_end and is_end_of_reply: + break + if has_is_end and (not is_end_of_reply) and (not next_marker): + warnings.append("marker_missing_when_not_end") + paging_note = "marker_missing_when_not_end" + break + if (not has_is_end) and (not next_marker): + if rows: + warnings.append("marker_missing_stop") + paging_note = "marker_missing_stop" + break + + meta = { + "page_count": page_no, + "last_marker": next_marker, + "is_end_of_reply": is_end_of_reply, + "graceful_end_by_iterator_error": graceful_end_by_iterator_error, + "paging_note": paging_note, + "warnings": warnings, + } + return pages, meta + + def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = "manual") -> UmeSyncJob: job = _build_sync_job("inventory", trigger_mode) db.add(job) db.flush() pulled = inserted = updated = 0 try: - ne_rows, _ = client.get_network_elements() + limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) + limit_max = max(1, limit_max) + page_size = int(getattr(settings, "ume_marker_page_limit", getattr(settings, "ume_page_size", 1000)) or 1000) + 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)) + + pages, _ = _collect_marker_pages( + lambda marker: client.get_network_elements(limit=page_size, marker=marker), + max_pages=max_pages, + iterator_500_as_end=False, + ) + ne_rows = [row for page in pages for row in page] now = _utc_now_naive() pulled = len(ne_rows) for row in ne_rows: @@ -198,6 +279,11 @@ def _sync_alarms_common( pulled = inserted = updated = 0 paging_mode = "marker" paging_note = "" + page_no = 0 + next_marker = "" + is_end_of_reply = False + graceful_end_by_iterator_error = False + warnings: list[str] = [] try: limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = max(1, limit_max) @@ -205,13 +291,6 @@ 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)) - page_no = 0 - next_marker = "" - last_page_signature = "" - is_end_of_reply = False - graceful_end_by_iterator_error = False - warnings: list[str] = [] - def upsert_alarm(alarm: dict[str, Any]) -> None: nonlocal inserted, updated key = _alarm_key(alarm) @@ -245,50 +324,21 @@ def _sync_alarms_common( existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str) iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True)) - - while True: - page_no += 1 - if page_no > max_pages: - raise RuntimeError(f"ume_alarms_pagination_exceeded:max_pages={max_pages}") - - try: - rows, diag = client.get_alarms( - is_uncleared=is_uncleared, - limit=page_size, - marker=(next_marker or None), - ) - except Exception as exc: - msg = str(exc or "") - low = msg.lower() - if iterator_500_as_end and pulled > 0 and "ume_request_failed:500" in low and "iterator" in low and "null" in low: - graceful_end_by_iterator_error = True - paging_note = msg[:240] - break - raise - + pages, meta = _collect_marker_pages( + lambda marker: client.get_alarms(is_uncleared=is_uncleared, limit=page_size, marker=marker), + max_pages=max_pages, + iterator_500_as_end=iterator_500_as_end, + ) + for rows in pages: pulled += len(rows) for alarm in rows: upsert_alarm(alarm) - - # Protection: if server keeps returning same page, stop to avoid infinite loop. - cur_sig = "|".join(sorted(_alarm_key(x) for x in rows if isinstance(x, dict))) - if cur_sig and cur_sig == last_page_signature: - warnings.append("duplicate_page_detected") - paging_note = "duplicate_page_detected" - break - last_page_signature = cur_sig - - is_end_of_reply = bool(diag.is_end_of_reply) if diag.is_end_of_reply is not None else False - next_marker = str(diag.marker or "").strip() - if is_end_of_reply: - break - # marker paging: if response header has no marker, treat as end of iteration. - # Some UME deployments omit marker when the first page already contains all rows. - if not next_marker: - if rows: - warnings.append("marker_missing_stop") - paging_note = "marker_missing_stop" - break + 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")) + graceful_end_by_iterator_error = bool(meta.get("graceful_end_by_iterator_error")) + 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) diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index a2a5299..44a53d1 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -138,7 +138,11 @@ class UmeSyncServiceTests(unittest.TestCase): def test_sync_inventory_upsert(self): class _C: - def get_network_elements(self): + def get_network_elements(self, *, limit=None, marker=None): + class _D: + marker = "" + is_end_of_reply = True + rows = [ { "ne-id": "NE-1", @@ -152,7 +156,7 @@ class UmeSyncServiceTests(unittest.TestCase): "vendor-name": "ZTE", }, ] - return rows, None + return rows, _D() job1 = sync_inventory_full(self.db, _C(), trigger_mode="manual") self.assertEqual(job1.status, "done") @@ -169,6 +173,36 @@ class UmeSyncServiceTests(unittest.TestCase): self.assertEqual(ne.host_name, "host-1") self.assertEqual(ne.hardware_version, "V1") + def test_sync_inventory_marker_pagination(self): + class _C: + def get_network_elements(self, *, limit=None, marker=None): + class _D: + def __init__(self, mk: str, end: bool): + self.marker = mk + self.is_end_of_reply = end + + if not marker: + return ( + [ + {"ne-id": "NE-1", "name": "ne1"}, + ], + _D("NM2", False), + ) + if marker == "NM2": + return ( + [ + {"ne-id": "NE-2", "name": "ne2"}, + ], + _D("NM3", True), + ) + return ([], _D("", True)) + + job = sync_inventory_full(self.db, _C(), trigger_mode="manual") + self.assertEqual(job.status, "done") + self.assertEqual(job.pulled_count, 2) + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-1")) + self.assertIsNotNone(self.db.get(UmeInventoryNE, "NE-2")) + def test_sync_current_alarms_upsert(self): class _C: def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): @@ -307,6 +341,32 @@ class UmeSyncServiceTests(unittest.TestCase): self.assertEqual(c.calls, 1) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-1")) + def test_sync_current_alarms_stop_when_not_end_but_marker_missing(self): + class _C: + def __init__(self): + self.calls = 0 + + def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): + self.calls += 1 + + class _D: + marker = "" + is_end_of_reply = False + + return ( + [ + {"alarmKey": "AK-9", "ne-id": "NE-9", "perceivedSeverity": "major", "isCleared": "false"}, + ], + _D(), + ) + + c = _C() + job, _ = sync_alarms_current(self.db, c, trigger_mode="manual") + self.assertEqual(job.status, "done") + self.assertEqual(job.pulled_count, 1) + self.assertEqual(c.calls, 1) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-9")) + def test_derive_ne_id_from_alarmkey_formats(self): alarm_hash = {"alarmkey": "00ceb960-1b62-478e-8303-0935ffea1d28#99010"} alarm_csv = {"alarmkey": "00ceb960-1b62-478e-8303-0935ffea1d28, 4237, 79"}