feat(UME): 切换marker分页并简化告警表字段

告警同步改为按响应头 marker 连续迭代,支持 iterator is null 的500尾页兜底,同时移除告警表中的 ne_name/user_label 冗余字段,统一在查询层与网元表 join 补齐展示信息。

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-08 18:58:24 +08:00
parent 3ea4f08bea
commit 267eddb273
8 changed files with 181 additions and 111 deletions

1
.gitignore vendored
View file

@ -6,3 +6,4 @@ __pycache__/
build/ build/
dist/ dist/
scripts/.run/ scripts/.run/
ume/

View file

@ -30,6 +30,9 @@ class Settings(BaseSettings):
ume_max_pages: int = 2000 ume_max_pages: int = 2000
ume_limit_max: int = 5000 ume_limit_max: int = 5000
ume_limit_only_page_size: 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_auth_header: str = "accessToken"
ume_content_type: str = "application/yang-data+json;charset=UTF-8" ume_content_type: str = "application/yang-data+json;charset=UTF-8"
ume_token_ttl_s: int = 1800 ume_token_ttl_s: int = 1800

View file

@ -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 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 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_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: except Exception:
pass pass
try: try:
@ -484,8 +489,6 @@ def ume_list_alarms(
stmt = stmt.filter( stmt = stmt.filter(
UmeAlarmCurrent.alarm_key.contains(kw) UmeAlarmCurrent.alarm_key.contains(kw)
| UmeAlarmCurrent.object_name.contains(kw) | UmeAlarmCurrent.object_name.contains(kw)
| UmeAlarmCurrent.ne_name.contains(kw)
| UmeAlarmCurrent.user_label.contains(kw)
| UmeAlarmCurrent.native_probable_cause.contains(kw) | UmeAlarmCurrent.native_probable_cause.contains(kw)
| UmeInventoryNE.ne_name.contains(kw) | UmeInventoryNE.ne_name.contains(kw)
| UmeInventoryNE.user_label.contains(kw) | UmeInventoryNE.user_label.contains(kw)
@ -497,8 +500,8 @@ def ume_list_alarms(
{ {
"alarm_key": str(alarm.alarm_key or ""), "alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(alarm.ne_id or ""), "ne_id": str(alarm.ne_id or ""),
"ne_name": str(alarm.ne_name or (ne.ne_name if ne else "") or ""), "ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str(alarm.user_label or (ne.user_label if ne else "") or ""), "user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(alarm.object_name or ""), "object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""), "event_type": str(alarm.event_type or ""),
"native_probable_cause": str(alarm.native_probable_cause or ""), "native_probable_cause": str(alarm.native_probable_cause or ""),
@ -514,9 +517,11 @@ def ume_list_alarms(
@app.get("/v1/ume/alarms/aggregate") @app.get("/v1/ume/alarms/aggregate")
def ume_alarms_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]: def ume_alarms_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]:
rows = db.query(UmeAlarmCurrent).all() rows = db.query(UmeAlarmCurrent, UmeInventoryNE).outerjoin(
by_severity = _aggregate_rows(rows, lambda x: x.perceived_severity) UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id
by_ne = _aggregate_rows(rows, lambda x: x.user_label or x.ne_name or x.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} return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne}
@ -543,8 +548,6 @@ def ume_list_alarms_history(
stmt = stmt.filter( stmt = stmt.filter(
UmeAlarmHistory.alarm_key.contains(kw) UmeAlarmHistory.alarm_key.contains(kw)
| UmeAlarmHistory.object_name.contains(kw) | UmeAlarmHistory.object_name.contains(kw)
| UmeAlarmHistory.ne_name.contains(kw)
| UmeAlarmHistory.user_label.contains(kw)
| UmeAlarmHistory.native_probable_cause.contains(kw) | UmeAlarmHistory.native_probable_cause.contains(kw)
| UmeInventoryNE.ne_name.contains(kw) | UmeInventoryNE.ne_name.contains(kw)
| UmeInventoryNE.user_label.contains(kw) | UmeInventoryNE.user_label.contains(kw)
@ -562,8 +565,8 @@ def ume_list_alarms_history(
{ {
"alarm_key": str(alarm.alarm_key or ""), "alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(alarm.ne_id or ""), "ne_id": str(alarm.ne_id or ""),
"ne_name": str(alarm.ne_name or (ne.ne_name if ne else "") or ""), "ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str(alarm.user_label or (ne.user_label if ne else "") or ""), "user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(alarm.object_name or ""), "object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""), "event_type": str(alarm.event_type or ""),
"native_probable_cause": str(alarm.native_probable_cause 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") @app.get("/v1/ume/alarms/history/aggregate")
def ume_alarms_history_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]: def ume_alarms_history_aggregate(db: Session = Depends(get_db)) -> dict[str, Any]:
rows = db.query(UmeAlarmHistory).all() rows = db.query(UmeAlarmHistory, UmeInventoryNE).outerjoin(
by_severity = _aggregate_rows(rows, lambda x: x.perceived_severity) UmeInventoryNE, UmeAlarmHistory.ne_id == UmeInventoryNE.ne_id
by_ne = _aggregate_rows(rows, lambda x: x.user_label or x.ne_name or x.ne_id) ).all()
by_date = _aggregate_rows(rows, lambda x: str(x.time_created or "")[:10]) 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} return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne, "by_date": by_date}

View file

@ -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)}]} return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]}
if name == "umeListCurrentAlarms": 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() severity = str(args.get("severity") or "").strip()
is_cleared = str(args.get("is_cleared") or "").strip() is_cleared = str(args.get("is_cleared") or "").strip()
ne_id = str(args.get("ne_id") 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( stmt = stmt.filter(
UmeAlarmCurrent.alarm_key.contains(keyword) UmeAlarmCurrent.alarm_key.contains(keyword)
| UmeAlarmCurrent.object_name.contains(keyword) | UmeAlarmCurrent.object_name.contains(keyword)
| UmeAlarmCurrent.ne_name.contains(keyword) | UmeInventoryNE.ne_name.contains(keyword)
| UmeAlarmCurrent.user_label.contains(keyword) | UmeInventoryNE.user_label.contains(keyword)
| UmeInventoryNE.ip_address.contains(keyword)
) )
page = max(1, int(args.get("page") or 1)) page = max(1, int(args.get("page") or 1))
page_size = min(500, max(1, int(args.get("page_size") or 50))) 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, "page_size": page_size,
"items": [ "items": [
{ {
"alarm_key": str(x.alarm_key or ""), "alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(x.ne_id or ""), "ne_id": str(alarm.ne_id or ""),
"ne_name": str(x.ne_name or ""), "ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str(x.user_label or ""), "user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(x.object_name or ""), "object_name": str(alarm.object_name or ""),
"perceived_severity": str(x.perceived_severity or ""), "perceived_severity": str(alarm.perceived_severity or ""),
"is_cleared": str(x.is_cleared or ""), "is_cleared": str(alarm.is_cleared or ""),
"time_created": str(x.time_created 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)}]} return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]}
if name == "umeListHistoryAlarms": 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() severity = str(args.get("severity") or "").strip()
ne_id = str(args.get("ne_id") or "").strip() ne_id = str(args.get("ne_id") or "").strip()
keyword = str(args.get("keyword") 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( stmt = stmt.filter(
UmeAlarmHistory.alarm_key.contains(keyword) UmeAlarmHistory.alarm_key.contains(keyword)
| UmeAlarmHistory.object_name.contains(keyword) | UmeAlarmHistory.object_name.contains(keyword)
| UmeAlarmHistory.ne_name.contains(keyword) | UmeInventoryNE.ne_name.contains(keyword)
| UmeAlarmHistory.user_label.contains(keyword) | UmeInventoryNE.user_label.contains(keyword)
| UmeInventoryNE.ip_address.contains(keyword)
) )
page = max(1, int(args.get("page") or 1)) page = max(1, int(args.get("page") or 1))
page_size = min(500, max(1, int(args.get("page_size") or 50))) 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, "page_size": page_size,
"items": [ "items": [
{ {
"alarm_key": str(x.alarm_key or ""), "alarm_key": str(alarm.alarm_key or ""),
"ne_id": str(x.ne_id or ""), "ne_id": str(alarm.ne_id or ""),
"ne_name": str(x.ne_name or ""), "ne_name": str((ne.ne_name if ne else "") or ""),
"user_label": str(x.user_label or ""), "user_label": str((ne.user_label if ne else "") or ""),
"object_name": str(x.object_name or ""), "object_name": str(alarm.object_name or ""),
"perceived_severity": str(x.perceived_severity or ""), "perceived_severity": str(alarm.perceived_severity or ""),
"is_cleared": str(x.is_cleared or ""), "is_cleared": str(alarm.is_cleared or ""),
"time_created": str(x.time_created 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)}]} 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 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 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_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: except Exception:
pass pass
for line in sys.stdin: for line in sys.stdin:

View file

@ -180,8 +180,6 @@ class UmeAlarmCurrent(Base):
alarm_key: Mapped[str] = mapped_column(String(256), primary_key=True) alarm_key: Mapped[str] = mapped_column(String(256), primary_key=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=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) object_name: Mapped[str] = mapped_column(String(512), default="", index=True)
event_type: Mapped[str] = mapped_column(String(128), default="") event_type: Mapped[str] = mapped_column(String(128), default="")
native_probable_cause: Mapped[str] = mapped_column(String(256), 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) alarm_key: Mapped[str] = mapped_column(String(256), primary_key=True)
ne_id: Mapped[str] = mapped_column(String(128), default="", index=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) object_name: Mapped[str] = mapped_column(String(512), default="", index=True)
event_type: Mapped[str] = mapped_column(String(128), default="") event_type: Mapped[str] = mapped_column(String(128), default="")
native_probable_cause: Mapped[str] = mapped_column(String(256), default="") native_probable_cause: Mapped[str] = mapped_column(String(256), default="")

View file

@ -31,6 +31,8 @@ class RequestDiagnostics:
latency_ms: int latency_ms: int
retry_count: int = 0 retry_count: int = 0
error_code: str = "" error_code: str = ""
marker: str = ""
is_end_of_reply: bool | None = None
class UMEClient: class UMEClient:
@ -361,6 +363,12 @@ class UMEClient:
self.login(force=True) self.login(force=True)
with self._client() as client: with self._client() as client:
resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True)) 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: if not resp.is_success:
diag = RequestDiagnostics( diag = RequestDiagnostics(
method=m, method=m,
@ -369,6 +377,8 @@ class UMEClient:
latency_ms=int((time() - t0) * 1000), latency_ms=int((time() - t0) * 1000),
retry_count=retry_count, retry_count=retry_count,
error_code=f"http_{int(resp.status_code)}", 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]}") raise RuntimeError(f"ume_request_failed:{resp.status_code}:{resp.text[:240]}")
data = _coerce_dict(resp.json()) data = _coerce_dict(resp.json())
@ -378,6 +388,8 @@ class UMEClient:
status_code=int(resp.status_code), status_code=int(resp.status_code),
latency_ms=int((time() - t0) * 1000), latency_ms=int((time() - t0) * 1000),
retry_count=retry_count, retry_count=retry_count,
marker=marker,
is_end_of_reply=is_end_of_reply,
) )
return data, diag return data, diag
except Exception as exc: except Exception as exc:
@ -446,7 +458,7 @@ class UMEClient:
*, *,
is_uncleared: bool, is_uncleared: bool,
limit: int | None = None, limit: int | None = None,
offset: int | None = None, marker: str | None = None,
) -> tuple[list[dict[str, Any]], RequestDiagnostics]: ) -> tuple[list[dict[str, Any]], RequestDiagnostics]:
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)
@ -456,8 +468,9 @@ class UMEClient:
"is-uncleared": "true" if is_uncleared else "false", "is-uncleared": "true" if is_uncleared else "false",
"limit": page_size, "limit": page_size,
} }
if offset is not None and int(offset) >= 0: marker_value = str(marker or "").strip()
params["offset"] = int(offset) if marker_value:
params["marker"] = marker_value
data, diag = self.request_json("GET", self.alarms_path, params=params) data, diag = self.request_json("GET", self.alarms_path, params=params)
rows = self._extract_named_list(data, ["alarm-list", "alarm"]) rows = self._extract_named_list(data, ["alarm-list", "alarm"])
return rows, diag return rows, diag

View file

@ -196,17 +196,21 @@ def _sync_alarms_common(
db.flush() db.flush()
now = _utc_now_naive() now = _utc_now_naive()
pulled = inserted = updated = 0 pulled = inserted = updated = 0
paging_mode = "offset" paging_mode = "marker"
paging_note = "" paging_note = ""
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)
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)) 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)) max_pages = max(1, min(max_pages, 20000))
offset = 0
page_no = 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: def upsert_alarm(alarm: dict[str, Any]) -> None:
nonlocal inserted, updated nonlocal inserted, updated
@ -228,8 +232,6 @@ def _sync_alarms_common(
else: else:
updated += 1 updated += 1
existing.ne_id = _derive_ne_id_from_alarm(alarm) 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.object_name = _s(_pick(alarm, "objectName", "object-name"))
existing.event_type = _s(_pick(alarm, "eventType", "event-type")) existing.event_type = _s(_pick(alarm, "eventType", "event-type"))
existing.native_probable_cause = _s(_pick(alarm, "nativeProbableCause", "native-probable-cause")) 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.last_seen_at = now
existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str) existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str)
def is_offset_unsupported_error(exc: Exception) -> bool: iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True))
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)
# Try offset pagination first (best-effort). If server rejects offset, fall back to single-page. while True:
try: page_no += 1
while True: if page_no > max_pages:
page_no += 1 raise RuntimeError(f"ume_alarms_pagination_exceeded:max_pages={max_pages}")
if page_no > max_pages:
raise RuntimeError(f"ume_alarms_pagination_exceeded:max_pages={max_pages}") try:
rows, _ = client.get_alarms(is_uncleared=is_uncleared, limit=page_size, offset=offset) rows, diag = client.get_alarms(
pulled += len(rows) is_uncleared=is_uncleared,
for alarm in rows: limit=page_size,
upsert_alarm(alarm) marker=(next_marker or None),
if len(rows) < page_size: )
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 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 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: if not is_uncleared:
expiry = now - timedelta(hours=48) expiry = now - timedelta(hours=48)
( (
@ -291,7 +300,17 @@ def _sync_alarms_common(
batch.status = "done" batch.status = "done"
batch.ended_at = _utc_now_naive() batch.ended_at = _utc_now_naive()
batch.raw_json = json.dumps( 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, ensure_ascii=False,
) )
@ -315,6 +334,11 @@ def _sync_alarms_common(
"status": batch.status, "status": batch.status,
"paging_mode": paging_mode, "paging_mode": paging_mode,
"paging_note": paging_note, "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, ensure_ascii=False,
) )

View file

@ -13,11 +13,12 @@ from netx_api.ume_sync_service import _derive_ne_id_from_alarm, sync_alarms_curr
class _FakeResponse: 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.status_code = int(status_code)
self._payload = payload self._payload = payload
self.text = str(payload) self.text = str(payload)
self.is_success = 200 <= self.status_code < 300 self.is_success = 200 <= self.status_code < 300
self.headers = headers or {}
def json(self): def json(self):
return self._payload return self._payload
@ -170,7 +171,7 @@ class UmeSyncServiceTests(unittest.TestCase):
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, offset=None): def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None):
rows = [ rows = [
{ {
"alarmKey": "AK-1", "alarmKey": "AK-1",
@ -185,7 +186,11 @@ class UmeSyncServiceTests(unittest.TestCase):
"timeCreated": "2026-01-01T00:00:00Z", "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") job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual")
self.assertEqual(job1.status, "done") self.assertEqual(job1.status, "done")
@ -199,27 +204,30 @@ class UmeSyncServiceTests(unittest.TestCase):
self.assertIsNotNone(row) self.assertIsNotNone(row)
self.assertEqual(row.perceived_severity, "critical") self.assertEqual(row.perceived_severity, "critical")
def test_sync_current_alarms_pagination(self): def test_sync_current_alarms_marker_pagination(self):
class _C: class _C:
def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None):
# Return 2 pages with page_size=2 then stop. class _D:
off = int(offset or 0) def __init__(self, mk: str, end: bool):
if off == 0: self.marker = mk
self.is_end_of_reply = end
if not marker:
return ( return (
[ [
{"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"}, {"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"},
{"alarmKey": "AK-2", "ne-id": "NE-2", "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 ( return (
[ [
{"alarmKey": "AK-3", "ne-id": "NE-3", "perceivedSeverity": "minor", "isCleared": "false"}, {"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 from netx_api import ume_sync_service as svc
old_page_size = svc.settings.ume_page_size old_page_size = svc.settings.ume_page_size
@ -235,29 +243,39 @@ class UmeSyncServiceTests(unittest.TestCase):
self.assertEqual(job.pulled_count, 3) self.assertEqual(job.pulled_count, 3)
self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-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: class _C:
def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): def get_alarms(self, *, is_uncleared: bool, limit=None, marker=None):
if offset is not None: class _D:
raise RuntimeError("ume_request_failed:400:unknown_param_offset") def __init__(self, mk: str, end: bool):
return ( self.marker = mk
[ self.is_end_of_reply = end
{"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"},
], if not marker:
None, 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 from netx_api import ume_sync_service as svc
old_page_size = svc.settings.ume_page_size old_page_size = svc.settings.ume_marker_page_limit
old_max_pages = svc.settings.ume_max_pages old_max_pages = svc.settings.ume_marker_max_pages
svc.settings.ume_page_size = 1000 old_500_as_end = svc.settings.ume_iterator_500_as_end
svc.settings.ume_max_pages = 10 svc.settings.ume_marker_page_limit = 1000
svc.settings.ume_marker_max_pages = 10
svc.settings.ume_iterator_500_as_end = True
try: try:
job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") job, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual")
finally: finally:
svc.settings.ume_page_size = old_page_size svc.settings.ume_marker_page_limit = old_page_size
svc.settings.ume_max_pages = old_max_pages 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.status, "done")
self.assertEqual(job.pulled_count, 1) self.assertEqual(job.pulled_count, 1)