From 267eddb27365317907dc3ccc03486b580d5b5a7a Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 8 May 2026 18:58:24 +0800 Subject: [PATCH] =?UTF-8?q?feat(UME):=20=E5=88=87=E6=8D=A2marker=E5=88=86?= =?UTF-8?q?=E9=A1=B5=E5=B9=B6=E7=AE=80=E5=8C=96=E5=91=8A=E8=AD=A6=E8=A1=A8?= =?UTF-8?q?=E5=AD=97=E6=AE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 告警同步改为按响应头 marker 连续迭代,支持 iterator is null 的500尾页兜底,同时移除告警表中的 ne_name/user_label 冗余字段,统一在查询层与网元表 join 补齐展示信息。 Co-authored-by: Cursor --- .gitignore | 1 + netx_api/config.py | 3 ++ netx_api/main.py | 35 ++++++------ netx_api/mcp_server.py | 58 +++++++++++--------- netx_api/models.py | 4 -- netx_api/ume_client.py | 19 +++++-- netx_api/ume_sync_service.py | 100 ++++++++++++++++++++++------------- tests/test_ume_sync.py | 72 +++++++++++++++---------- 8 files changed, 181 insertions(+), 111 deletions(-) diff --git a/.gitignore b/.gitignore index 540a6f9..2fe42e2 100644 --- a/.gitignore +++ b/.gitignore @@ -6,3 +6,4 @@ __pycache__/ build/ dist/ scripts/.run/ +ume/ diff --git a/netx_api/config.py b/netx_api/config.py index faf2662..dcb2b0d 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -30,6 +30,9 @@ class Settings(BaseSettings): ume_max_pages: int = 2000 ume_limit_max: int = 5000 ume_limit_only_page_size: int = 5000 + ume_marker_page_limit: int = 1000 + ume_marker_max_pages: int = 2000 + ume_iterator_500_as_end: bool = True ume_auth_header: str = "accessToken" ume_content_type: str = "application/yang-data+json;charset=UTF-8" ume_token_ttl_s: int = 1800 diff --git a/netx_api/main.py b/netx_api/main.py index 3da73c5..ddd04ac 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -220,6 +220,11 @@ def on_startup() -> None: conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS net_mask VARCHAR(128) DEFAULT ''") conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS create_time VARCHAR(64) DEFAULT ''") conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS creator VARCHAR(128) DEFAULT ''") + # Simplify alarm tables: display fields come from runtime join with inventory table. + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS user_label") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS user_label") except Exception: pass try: @@ -484,8 +489,6 @@ def ume_list_alarms( stmt = stmt.filter( UmeAlarmCurrent.alarm_key.contains(kw) | UmeAlarmCurrent.object_name.contains(kw) - | UmeAlarmCurrent.ne_name.contains(kw) - | UmeAlarmCurrent.user_label.contains(kw) | UmeAlarmCurrent.native_probable_cause.contains(kw) | UmeInventoryNE.ne_name.contains(kw) | UmeInventoryNE.user_label.contains(kw) @@ -497,8 +500,8 @@ def ume_list_alarms( { "alarm_key": str(alarm.alarm_key or ""), "ne_id": str(alarm.ne_id or ""), - "ne_name": str(alarm.ne_name or (ne.ne_name if ne else "") or ""), - "user_label": str(alarm.user_label or (ne.user_label if ne else "") or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), "object_name": str(alarm.object_name or ""), "event_type": str(alarm.event_type or ""), "native_probable_cause": str(alarm.native_probable_cause or ""), @@ -514,9 +517,11 @@ def ume_list_alarms( @app.get("/v1/ume/alarms/aggregate") def ume_alarms_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]: - rows = db.query(UmeAlarmCurrent).all() - by_severity = _aggregate_rows(rows, lambda x: x.perceived_severity) - by_ne = _aggregate_rows(rows, lambda x: x.user_label or x.ne_name or x.ne_id) + rows = db.query(UmeAlarmCurrent, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id + ).all() + by_severity = _aggregate_rows(rows, lambda x: x[0].perceived_severity) + by_ne = _aggregate_rows(rows, lambda x: (x[1].user_label if x[1] else "") or (x[1].ne_name if x[1] else "") or x[0].ne_id) return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne} @@ -543,8 +548,6 @@ def ume_list_alarms_history( stmt = stmt.filter( UmeAlarmHistory.alarm_key.contains(kw) | UmeAlarmHistory.object_name.contains(kw) - | UmeAlarmHistory.ne_name.contains(kw) - | UmeAlarmHistory.user_label.contains(kw) | UmeAlarmHistory.native_probable_cause.contains(kw) | UmeInventoryNE.ne_name.contains(kw) | UmeInventoryNE.user_label.contains(kw) @@ -562,8 +565,8 @@ def ume_list_alarms_history( { "alarm_key": str(alarm.alarm_key or ""), "ne_id": str(alarm.ne_id or ""), - "ne_name": str(alarm.ne_name or (ne.ne_name if ne else "") or ""), - "user_label": str(alarm.user_label or (ne.user_label if ne else "") or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), "object_name": str(alarm.object_name or ""), "event_type": str(alarm.event_type or ""), "native_probable_cause": str(alarm.native_probable_cause or ""), @@ -579,10 +582,12 @@ def ume_list_alarms_history( @app.get("/v1/ume/alarms/history/aggregate") def ume_alarms_history_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]: - rows = db.query(UmeAlarmHistory).all() - by_severity = _aggregate_rows(rows, lambda x: x.perceived_severity) - by_ne = _aggregate_rows(rows, lambda x: x.user_label or x.ne_name or x.ne_id) - by_date = _aggregate_rows(rows, lambda x: str(x.time_created or "")[:10]) + rows = db.query(UmeAlarmHistory, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmHistory.ne_id == UmeInventoryNE.ne_id + ).all() + by_severity = _aggregate_rows(rows, lambda x: x[0].perceived_severity) + by_ne = _aggregate_rows(rows, lambda x: (x[1].user_label if x[1] else "") or (x[1].ne_name if x[1] else "") or x[0].ne_id) + by_date = _aggregate_rows(rows, lambda x: str(x[0].time_created or "")[:10]) return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne, "by_date": by_date} diff --git a/netx_api/mcp_server.py b/netx_api/mcp_server.py index 5c36a6c..a08ecef 100644 --- a/netx_api/mcp_server.py +++ b/netx_api/mcp_server.py @@ -346,7 +346,9 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: } return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} if name == "umeListCurrentAlarms": - stmt = db.query(UmeAlarmCurrent) + stmt = db.query(UmeAlarmCurrent, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id + ) severity = str(args.get("severity") or "").strip() is_cleared = str(args.get("is_cleared") or "").strip() ne_id = str(args.get("ne_id") or "").strip() @@ -361,8 +363,9 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: stmt = stmt.filter( UmeAlarmCurrent.alarm_key.contains(keyword) | UmeAlarmCurrent.object_name.contains(keyword) - | UmeAlarmCurrent.ne_name.contains(keyword) - | UmeAlarmCurrent.user_label.contains(keyword) + | UmeInventoryNE.ne_name.contains(keyword) + | UmeInventoryNE.user_label.contains(keyword) + | UmeInventoryNE.ip_address.contains(keyword) ) page = max(1, int(args.get("page") or 1)) page_size = min(500, max(1, int(args.get("page_size") or 50))) @@ -374,21 +377,23 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: "page_size": page_size, "items": [ { - "alarm_key": str(x.alarm_key or ""), - "ne_id": str(x.ne_id or ""), - "ne_name": str(x.ne_name or ""), - "user_label": str(x.user_label or ""), - "object_name": str(x.object_name or ""), - "perceived_severity": str(x.perceived_severity or ""), - "is_cleared": str(x.is_cleared or ""), - "time_created": str(x.time_created or ""), + "alarm_key": str(alarm.alarm_key or ""), + "ne_id": str(alarm.ne_id or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), + "object_name": str(alarm.object_name or ""), + "perceived_severity": str(alarm.perceived_severity or ""), + "is_cleared": str(alarm.is_cleared or ""), + "time_created": str(alarm.time_created or ""), } - for x in rows + for alarm, ne in rows ], } return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} if name == "umeListHistoryAlarms": - stmt = db.query(UmeAlarmHistory) + stmt = db.query(UmeAlarmHistory, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmHistory.ne_id == UmeInventoryNE.ne_id + ) severity = str(args.get("severity") or "").strip() ne_id = str(args.get("ne_id") or "").strip() keyword = str(args.get("keyword") or "").strip() @@ -400,8 +405,9 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: stmt = stmt.filter( UmeAlarmHistory.alarm_key.contains(keyword) | UmeAlarmHistory.object_name.contains(keyword) - | UmeAlarmHistory.ne_name.contains(keyword) - | UmeAlarmHistory.user_label.contains(keyword) + | UmeInventoryNE.ne_name.contains(keyword) + | UmeInventoryNE.user_label.contains(keyword) + | UmeInventoryNE.ip_address.contains(keyword) ) page = max(1, int(args.get("page") or 1)) page_size = min(500, max(1, int(args.get("page_size") or 50))) @@ -413,16 +419,16 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: "page_size": page_size, "items": [ { - "alarm_key": str(x.alarm_key or ""), - "ne_id": str(x.ne_id or ""), - "ne_name": str(x.ne_name or ""), - "user_label": str(x.user_label or ""), - "object_name": str(x.object_name or ""), - "perceived_severity": str(x.perceived_severity or ""), - "is_cleared": str(x.is_cleared or ""), - "time_created": str(x.time_created or ""), + "alarm_key": str(alarm.alarm_key or ""), + "ne_id": str(alarm.ne_id or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), + "object_name": str(alarm.object_name or ""), + "perceived_severity": str(alarm.perceived_severity or ""), + "is_cleared": str(alarm.is_cleared or ""), + "time_created": str(alarm.time_created or ""), } - for x in rows + for alarm, ne in rows ], } return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} @@ -458,6 +464,10 @@ def main() -> None: conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS net_mask VARCHAR(128) DEFAULT ''") conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS create_time VARCHAR(64) DEFAULT ''") conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS creator VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS user_label") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS user_label") except Exception: pass for line in sys.stdin: diff --git a/netx_api/models.py b/netx_api/models.py index 263d9b3..7583511 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -180,8 +180,6 @@ class UmeAlarmCurrent(Base): alarm_key: Mapped[str] = mapped_column(String(256), primary_key=True) ne_id: Mapped[str] = mapped_column(String(128), default="", index=True) - ne_name: Mapped[str] = mapped_column(String(256), default="", index=True) - user_label: Mapped[str] = mapped_column(String(256), default="", index=True) object_name: Mapped[str] = mapped_column(String(512), default="", index=True) event_type: Mapped[str] = mapped_column(String(128), default="") native_probable_cause: Mapped[str] = mapped_column(String(256), default="") @@ -199,8 +197,6 @@ class UmeAlarmHistory(Base): alarm_key: Mapped[str] = mapped_column(String(256), primary_key=True) ne_id: Mapped[str] = mapped_column(String(128), default="", index=True) - ne_name: Mapped[str] = mapped_column(String(256), default="", index=True) - user_label: Mapped[str] = mapped_column(String(256), default="", index=True) object_name: Mapped[str] = mapped_column(String(512), default="", index=True) event_type: Mapped[str] = mapped_column(String(128), default="") native_probable_cause: Mapped[str] = mapped_column(String(256), default="") diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py index af564ab..ef8c512 100644 --- a/netx_api/ume_client.py +++ b/netx_api/ume_client.py @@ -31,6 +31,8 @@ class RequestDiagnostics: latency_ms: int retry_count: int = 0 error_code: str = "" + marker: str = "" + is_end_of_reply: bool | None = None class UMEClient: @@ -361,6 +363,12 @@ class UMEClient: self.login(force=True) with self._client() 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() + is_end_of_reply: bool | None = None + if is_end_raw in {"true", "false"}: + is_end_of_reply = is_end_raw == "true" + if not resp.is_success: diag = RequestDiagnostics( method=m, @@ -369,6 +377,8 @@ class UMEClient: latency_ms=int((time() - t0) * 1000), retry_count=retry_count, error_code=f"http_{int(resp.status_code)}", + marker=marker, + is_end_of_reply=is_end_of_reply, ) raise RuntimeError(f"ume_request_failed:{resp.status_code}:{resp.text[:240]}") data = _coerce_dict(resp.json()) @@ -378,6 +388,8 @@ class UMEClient: status_code=int(resp.status_code), latency_ms=int((time() - t0) * 1000), retry_count=retry_count, + marker=marker, + is_end_of_reply=is_end_of_reply, ) return data, diag except Exception as exc: @@ -446,7 +458,7 @@ class UMEClient: *, is_uncleared: bool, limit: int | None = None, - offset: 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) @@ -456,8 +468,9 @@ class UMEClient: "is-uncleared": "true" if is_uncleared else "false", "limit": page_size, } - if offset is not None and int(offset) >= 0: - params["offset"] = int(offset) + marker_value = str(marker or "").strip() + if marker_value: + params["marker"] = marker_value data, diag = self.request_json("GET", self.alarms_path, params=params) rows = self._extract_named_list(data, ["alarm-list", "alarm"]) return rows, diag diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index 97ef456..37f7c5a 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -196,17 +196,21 @@ def _sync_alarms_common( db.flush() now = _utc_now_naive() pulled = inserted = updated = 0 - paging_mode = "offset" + paging_mode = "marker" paging_note = "" try: limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) limit_max = max(1, limit_max) - page_size = int(getattr(settings, "ume_page_size", 1000) or 1000) + 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_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)) - offset = 0 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 @@ -228,8 +232,6 @@ def _sync_alarms_common( else: updated += 1 existing.ne_id = _derive_ne_id_from_alarm(alarm) - existing.ne_name = _s(_pick(alarm, "ne-name", "neName", "ne_name")) - existing.user_label = _s(_pick(alarm, "user-label", "userLabel", "user_label")) existing.object_name = _s(_pick(alarm, "objectName", "object-name")) existing.event_type = _s(_pick(alarm, "eventType", "event-type")) existing.native_probable_cause = _s(_pick(alarm, "nativeProbableCause", "native-probable-cause")) @@ -242,41 +244,48 @@ def _sync_alarms_common( existing.last_seen_at = now existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str) - def is_offset_unsupported_error(exc: Exception) -> bool: - msg = str(exc or "") - if "ume_request_failed:400" not in msg: - return False - low = msg.lower() - return ("offset" in low) or ("unknown" in low and "param" in low) or ("illegal" in low and "param" in low) + iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True)) - # Try offset pagination first (best-effort). If server rejects offset, fall back to single-page. - try: - while True: - page_no += 1 - if page_no > max_pages: - raise RuntimeError(f"ume_alarms_pagination_exceeded:max_pages={max_pages}") - rows, _ = client.get_alarms(is_uncleared=is_uncleared, limit=page_size, offset=offset) - pulled += len(rows) - for alarm in rows: - upsert_alarm(alarm) - if len(rows) < page_size: + 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 - offset += page_size - except Exception as exc: - if is_offset_unsupported_error(exc): - paging_mode = "limit_only" - paging_note = str(exc)[:200] - limit_only_size = int(getattr(settings, "ume_limit_only_page_size", limit_max) or limit_max) - limit_only_size = max(1, min(limit_only_size, limit_max)) - # Re-run as single page without offset param. - pulled = inserted = updated = 0 - rows, _ = client.get_alarms(is_uncleared=is_uncleared, limit=limit_only_size, offset=None) - pulled = len(rows) - for alarm in rows: - upsert_alarm(alarm) - else: raise + 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: no marker and empty data means no next page. + if not next_marker and not rows: + break + if not is_uncleared: expiry = now - timedelta(hours=48) ( @@ -291,7 +300,17 @@ def _sync_alarms_common( batch.status = "done" batch.ended_at = _utc_now_naive() batch.raw_json = json.dumps( - {"pulled": pulled, "inserted": inserted, "updated": updated, "paging_mode": paging_mode}, + { + "pulled": pulled, + "inserted": inserted, + "updated": updated, + "paging_mode": paging_mode, + "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, + "warnings": warnings, + }, ensure_ascii=False, ) @@ -315,6 +334,11 @@ def _sync_alarms_common( "status": batch.status, "paging_mode": paging_mode, "paging_note": paging_note, + "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, + "warnings": warnings, }, ensure_ascii=False, ) diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py index 9674232..f538617 100644 --- a/tests/test_ume_sync.py +++ b/tests/test_ume_sync.py @@ -13,11 +13,12 @@ from netx_api.ume_sync_service import _derive_ne_id_from_alarm, sync_alarms_curr class _FakeResponse: - def __init__(self, status_code: int, payload: dict): + def __init__(self, status_code: int, payload: dict, headers: dict | None = None): self.status_code = int(status_code) self._payload = payload self.text = str(payload) self.is_success = 200 <= self.status_code < 300 + self.headers = headers or {} def json(self): return self._payload @@ -170,7 +171,7 @@ class UmeSyncServiceTests(unittest.TestCase): def test_sync_current_alarms_upsert(self): class _C: - def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): + def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None): rows = [ { "alarmKey": "AK-1", @@ -185,7 +186,11 @@ class UmeSyncServiceTests(unittest.TestCase): "timeCreated": "2026-01-01T00:00:00Z", } ] - return rows, None + class _D: + marker = "" + is_end_of_reply = True + + return rows, _D() job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") self.assertEqual(job1.status, "done") @@ -199,27 +204,30 @@ class UmeSyncServiceTests(unittest.TestCase): self.assertIsNotNone(row) self.assertEqual(row.perceived_severity, "critical") - def test_sync_current_alarms_pagination(self): + def test_sync_current_alarms_marker_pagination(self): class _C: - def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): - # Return 2 pages with page_size=2 then stop. - off = int(offset or 0) - if off == 0: + def get_alarms(self, *, is_uncleared: bool, 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 ( [ {"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"}, {"alarmKey": "AK-2", "ne-id": "NE-2", "perceivedSeverity": "major", "isCleared": "false"}, ], - None, + _D("M2", False), ) - if off == 2: + if marker == "M2": return ( [ {"alarmKey": "AK-3", "ne-id": "NE-3", "perceivedSeverity": "minor", "isCleared": "false"}, ], - None, + _D("M3", True), ) - return ([], None) + return ([], _D("", True)) from netx_api import ume_sync_service as svc old_page_size = svc.settings.ume_page_size @@ -235,29 +243,39 @@ class UmeSyncServiceTests(unittest.TestCase): self.assertEqual(job.pulled_count, 3) self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-3")) - def test_sync_current_alarms_offset_unsupported_fallback(self): + def test_sync_current_alarms_iterator_500_as_end(self): class _C: - def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): - if offset is not None: - raise RuntimeError("ume_request_failed:400:unknown_param_offset") - return ( - [ - {"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"}, - ], - None, + def get_alarms(self, *, is_uncleared: bool, 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 ( + [ + {"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"}, + ], + _D("M2", False), + ) + raise RuntimeError( + "ume_request_failed:500:{\"error\":{\"errorCode\":\"500\",\"errorInfo\":\"iterator is null\"}}" ) from netx_api import ume_sync_service as svc - old_page_size = svc.settings.ume_page_size - old_max_pages = svc.settings.ume_max_pages - svc.settings.ume_page_size = 1000 - svc.settings.ume_max_pages = 10 + old_page_size = svc.settings.ume_marker_page_limit + old_max_pages = svc.settings.ume_marker_max_pages + old_500_as_end = svc.settings.ume_iterator_500_as_end + svc.settings.ume_marker_page_limit = 1000 + svc.settings.ume_marker_max_pages = 10 + svc.settings.ume_iterator_500_as_end = True try: job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") finally: - svc.settings.ume_page_size = old_page_size - svc.settings.ume_max_pages = old_max_pages + svc.settings.ume_marker_page_limit = old_page_size + svc.settings.ume_marker_max_pages = old_max_pages + svc.settings.ume_iterator_500_as_end = old_500_as_end self.assertEqual(job.status, "done") self.assertEqual(job.pulled_count, 1)