diff --git a/.env.example b/.env.example index 630bd79..9b71072 100644 --- a/.env.example +++ b/.env.example @@ -7,3 +7,15 @@ NETX_PARSER_CONFIG=netx_api/config/parsers/zte_alarm_monitor_v1.yaml NETX_OCLAW_ANALYZE_URL=http://127.0.0.1:8787/admin/api/ops-ai/analyze-sync NETX_OCLAW_ANALYZE_TOKEN= NETX_OCLAW_HEALTH_URL=http://127.0.0.1:8787/admin/api/ops-ai/health +NETX_UME_BASE_URL=https://10.227.157.143:18014 +NETX_UME_USERNAME= +NETX_UME_PASSWORD= +NETX_UME_AUTH_HEADER=accessToken +NETX_UME_CONTENT_TYPE=application/yang-data+json;charset=UTF-8 +NETX_UME_TOKEN_TTL_S=1800 +NETX_UME_TOKEN_REFRESH_SKEW_S=60 +NETX_UME_TOKEN_PATH=/restconf/operations/zte-security:oauth_token +NETX_UME_TOKEN_HANDSHAKE_PATH=/restconf/operations/zte-security:oauth_handshake +NETX_UME_TOKEN_LOGOUT_PATH=/restconf/operations/zte-security:oauth_token +NETX_UME_NE_PATH=/restconf/data/zte-resources-module:network-elements +NETX_UME_ALARMS_PATH=/restconf/data/zte-alarms:alarms/alarm-list diff --git a/netx_api/config.py b/netx_api/config.py index 67ef8a1..faf2662 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -20,6 +20,30 @@ class Settings(BaseSettings): oclaw_connect_timeout_sec: float = 15.0 oclaw_analyze_read_timeout_sec: float = 180.0 oclaw_health_timeout_sec: float = 8.0 + # UME RESTCONF integration + ume_base_url: str = "" + ume_username: str = "" + ume_password: str = "" + ume_verify_tls: bool = False + ume_timeout_s: float = 20.0 + ume_page_size: int = 1000 + ume_max_pages: int = 2000 + ume_limit_max: int = 5000 + ume_limit_only_page_size: int = 5000 + ume_auth_header: str = "accessToken" + ume_content_type: str = "application/yang-data+json;charset=UTF-8" + ume_token_ttl_s: int = 1800 + ume_token_refresh_skew_s: int = 60 + ume_keepalive_enabled: bool = True + ume_keepalive_interval_s: int = 600 + ume_keepalive_renew_before_s: int = 900 + ume_token_path: str = "/restconf/operations/zte-security:oauth_token" + ume_token_handshake_path: str = "/restconf/operations/zte-security:oauth_handshake" + ume_token_logout_path: str = "/restconf/operations/zte-security:oauth_token" + ume_ne_path: str = "/restconf/data/zte-resources-module:network-elements" + ume_alarms_path: str = "/restconf/data/zte-alarms:alarms/alarm-list" + ume_sync_inventory_every_hours: int = 24 + ume_sync_alarms_history_every_hours: int = 24 settings = Settings() diff --git a/netx_api/main.py b/netx_api/main.py index 820d6f3..0a881bb 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -6,6 +6,7 @@ from datetime import datetime, timezone from io import StringIO import time import re +import threading from fastapi import Depends, FastAPI, File, HTTPException, Query, UploadFile from fastapi.responses import Response from sqlalchemy import text as sql_text @@ -17,9 +18,28 @@ from .ap_client import analyze_with_oclaw, health_with_oclaw from .config import settings from .db import Base, SessionLocal, engine from .importer import aggregate_alarms, import_alarm_excel, query_alarms -from .models import AiAnalyzeHistory, AlarmBatch, AlarmNorm, ImportErrorRow +from .models import ( + AiAnalyzeHistory, + AlarmBatch, + AlarmNorm, + ImportErrorRow, + UmeAlarmCurrent, + UmeAlarmHistory, + UmeInventoryNE, + UmeSyncJob, +) from .models import ImportJob from .parser_config import load_parser_config +from .ume_client import UMEClient +from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full +from .ume_token_store import ( + clear_shared_token, + load_shared_token, + release_refresh_lock, + save_shared_token, + try_acquire_refresh_lock, + wait_for_token_update, +) from .schemas import ( AlarmAggregateBucket, AlarmAggregateResponse, @@ -34,6 +54,14 @@ from .schemas import ( app = FastAPI(title="netx ops tool", version="0.1.0") parser_cfg = load_parser_config() +_UME_CLIENT_SINGLETON = UMEClient( + token_loader=lambda: load_shared_token(), + token_saver=lambda token, exp: save_shared_token(token, exp), + token_clearer=lambda: clear_shared_token(), + lock_acquirer=lambda: try_acquire_refresh_lock(), + lock_releaser=lambda: release_refresh_lock(), + token_waiter=lambda min_exp: wait_for_token_update(min_expires_at_epoch_s=float(min_exp)), +) _SQL_FORBIDDEN_RE = re.compile( r"\b(insert|update|delete|drop|alter|create|truncate|grant|revoke|call|copy|vacuum|analyze)\b", @@ -61,6 +89,47 @@ def get_db(): db.close() +def _parse_time(text: str | None) -> datetime | None: + s = str(text or "").strip() + if not s: + return None + s2 = s.replace("Z", "+00:00") + try: + dt = datetime.fromisoformat(s2) + return _ensure_utc(dt) + except Exception: + return None + + +def _aggregate_rows(items: list[Any], key_fn) -> list[dict[str, Any]]: + bucket: dict[str, int] = {} + for item in items: + key = str(key_fn(item) or "").strip() + if not key: + key = "unknown" + bucket[key] = int(bucket.get(key, 0)) + 1 + return [{"key": k, "count": v} for k, v in sorted(bucket.items(), key=lambda kv: kv[1], reverse=True)] + + +def _ume_client() -> UMEClient: + return _UME_CLIENT_SINGLETON + + +def _ume_error_kind(err: str) -> str: + low = str(err or "").lower() + if "401" in low or "403" in low or "password" in low or "auth" in low: + return "auth_failed" + if "timeout" in low: + return "timeout" + if "tls" in low or "certificate" in low or "ssl" in low: + return "tls_failed" + if "connect" in low or "name or service not known" in low: + return "connect_failed" + if "handshake" in low: + return "handshake_failed" + return "other" + + @app.post("/v1/sql/query") def sql_query(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict: """ @@ -124,6 +193,32 @@ def on_startup() -> None: ) conn.exec_driver_sql("ALTER TABLE alarms_norm ADD COLUMN IF NOT EXISTS intermittence_count INTEGER DEFAULT 0") conn.exec_driver_sql("ALTER TABLE alarms_norm ADD COLUMN IF NOT EXISTS me_level VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_token_cache ADD COLUMN IF NOT EXISTS lock_owner VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_token_cache ADD COLUMN IF NOT EXISTS lock_expires_at_epoch_s INTEGER DEFAULT 0") + except Exception: + pass + try: + if bool(getattr(settings, "ume_keepalive_enabled", True)): + interval_s = int(getattr(settings, "ume_keepalive_interval_s", 600) or 600) + interval_s = max(30, min(interval_s, 3600)) + renew_before_s = int(getattr(settings, "ume_keepalive_renew_before_s", 900) or 900) + renew_before_s = max(30, min(renew_before_s, 86400)) + + def _keepalive_loop() -> None: + # Best-effort keepalive: if token exists, periodically handshake to extend TTL. + while True: + try: + client = _ume_client() + st = client.token_status() + expires_in = int(st.get("expires_in_s") or 0) + if bool(st.get("has_token")) and expires_in > 0 and expires_in < renew_before_s: + client.renew_token() + except Exception: + pass + time.sleep(interval_s) + + t = threading.Thread(target=_keepalive_loop, name="ume-token-keepalive", daemon=True) + t.start() except Exception: pass @@ -133,6 +228,299 @@ def health() -> dict[str, str]: return {"status": "ok"} +@app.get("/v1/ume/token/status") +def ume_token_status() -> dict[str, Any]: + client = _ume_client() + st = client.token_status() + return {"ok": True, **st} + + +@app.post("/v1/ume/token/refresh") +def ume_token_refresh() -> dict[str, Any]: + client = _ume_client() + try: + before = client.token_status() + token = client.refresh_if_needed() + after = client.token_status() + return { + "ok": True, + "token": token, + "changed": bool(before.get("token_preview") != after.get("token_preview")), + **after, + } + except Exception as exc: + msg = str(exc)[:240] + return {"ok": False, "error_kind": _ume_error_kind(msg), "error": msg} + + +@app.post("/v1/ume/token/disconnect") +def ume_token_disconnect() -> dict[str, Any]: + client = _ume_client() + ok = bool(client.logout_token()) + st = client.token_status() + return {"ok": ok, **st} + + +@app.post("/v1/ume/sync") +def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict[str, Any]: + body = payload or {} + domains = body.get("domains") + if not isinstance(domains, list) or not domains: + domains = ["inventory", "alarms_current", "alarms_history"] + domain_set = {str(x).strip().lower() for x in domains if str(x).strip()} + trigger_mode = str(body.get("trigger_mode") or "manual").strip().lower() + if trigger_mode not in {"manual", "schedule"}: + trigger_mode = "manual" + + client = _ume_client() + out: dict[str, Any] = {"ok": True, "jobs": []} + try: + if "inventory" in domain_set: + job = sync_inventory_full(db, client, trigger_mode=trigger_mode) + out["jobs"].append( + { + "domain": "inventory", + "status": job.status, + "pulled_count": int(job.pulled_count or 0), + "inserted_count": int(job.inserted_count or 0), + "updated_count": int(job.updated_count or 0), + "error_message": str(job.error_message or ""), + } + ) + if "alarms" in domain_set or "alarms_current" in domain_set: + job, batch = sync_alarms_current(db, client, trigger_mode=trigger_mode) + out["jobs"].append( + { + "domain": "alarms_current", + "status": job.status, + "batch_id": str(batch.batch_id), + "pulled_count": int(job.pulled_count or 0), + "inserted_count": int(job.inserted_count or 0), + "updated_count": int(job.updated_count or 0), + "error_message": str(job.error_message or ""), + } + ) + if "alarms_history" in domain_set: + job, batch = sync_alarms_history_full(db, client, trigger_mode=trigger_mode) + out["jobs"].append( + { + "domain": "alarms_history", + "status": job.status, + "batch_id": str(batch.batch_id), + "pulled_count": int(job.pulled_count or 0), + "inserted_count": int(job.inserted_count or 0), + "updated_count": int(job.updated_count or 0), + "error_message": str(job.error_message or ""), + } + ) + except Exception as exc: + out["ok"] = False + out["error"] = str(exc)[:240] + return out + + +@app.get("/v1/ume/sync/status") +def ume_sync_status( + page: int = Query(default=1, ge=1), + page_size: int = Query(default=20, ge=1, le=200), + db: Session = Depends(get_db), +) -> dict[str, Any]: + q = db.query(UmeSyncJob) + total = int(q.count()) + rows = ( + q.order_by(UmeSyncJob.id.desc()) + .offset((int(page) - 1) * int(page_size)) + .limit(int(page_size)) + .all() + ) + items = [] + latest_by_domain: dict[str, dict[str, Any]] = {} + for r in rows: + item = { + "id": int(r.id), + "domain": str(r.domain or ""), + "status": str(r.status or ""), + "trigger_mode": str(r.trigger_mode or ""), + "pulled_count": int(r.pulled_count or 0), + "inserted_count": int(r.inserted_count or 0), + "updated_count": int(r.updated_count or 0), + "error_message": str(r.error_message or ""), + "started_at": (_ensure_utc(r.started_at) or datetime.now(timezone.utc)).isoformat(), + "ended_at": (_ensure_utc(r.ended_at).isoformat() if r.ended_at else None), + } + items.append(item) + if item["domain"] and item["domain"] not in latest_by_domain: + latest_by_domain[item["domain"]] = item + return {"total": total, "page": page, "page_size": page_size, "items": items, "latest_by_domain": latest_by_domain} + + +@app.get("/v1/ume/inventory/ne") +def ume_list_ne( + keyword: str | None = Query(default=None), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=50, ge=1, le=500), + db: Session = Depends(get_db), +) -> dict[str, Any]: + stmt = db.query(UmeInventoryNE) + kw = str(keyword or "").strip() + if kw: + stmt = stmt.filter( + UmeInventoryNE.ne_id.contains(kw) + | UmeInventoryNE.ne_name.contains(kw) + | UmeInventoryNE.user_label.contains(kw) + | UmeInventoryNE.ip_address.contains(kw) + ) + total = int(stmt.count()) + rows = stmt.order_by(UmeInventoryNE.ne_id.asc()).offset((page - 1) * page_size).limit(page_size).all() + items = [ + { + "ne_id": str(x.ne_id or ""), + "ne_name": str(x.ne_name or ""), + "user_label": str(x.user_label or ""), + "ip_address": str(x.ip_address or ""), + "ne_type": str(x.ne_type or ""), + "last_seen_at": (_ensure_utc(x.last_seen_at) or datetime.now(timezone.utc)).isoformat(), + } + for x in rows + ] + return {"total": total, "page": page, "page_size": page_size, "items": items} + + +@app.get("/v1/ume/inventory/ne/{ne_id}") +def ume_get_ne(ne_id: str, db: Session = Depends(get_db)) -> dict[str, Any]: + row = db.get(UmeInventoryNE, ne_id) + if not row: + raise HTTPException(status_code=404, detail="ume_ne_not_found") + return { + "ne_id": str(row.ne_id or ""), + "ne_name": str(row.ne_name or ""), + "user_label": str(row.user_label or ""), + "ip_address": str(row.ip_address or ""), + "ne_type": str(row.ne_type or ""), + "vendor": str(row.vendor or ""), + "source_type": str(row.source_type or ""), + "first_seen_at": (_ensure_utc(row.first_seen_at) or datetime.now(timezone.utc)).isoformat(), + "last_seen_at": (_ensure_utc(row.last_seen_at) or datetime.now(timezone.utc)).isoformat(), + "raw_json": str(row.raw_json or "{}"), + } + + +@app.get("/v1/ume/alarms") +def ume_list_alarms( + severity: str | None = Query(default=None), + is_cleared: str | None = Query(default=None), + ne_id: str | None = Query(default=None), + keyword: str | None = Query(default=None), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=50, ge=1, le=500), + db: Session = Depends(get_db), +) -> dict[str, Any]: + stmt = db.query(UmeAlarmCurrent) + if severity and str(severity).strip(): + stmt = stmt.filter(UmeAlarmCurrent.perceived_severity == str(severity).strip()) + if is_cleared and str(is_cleared).strip(): + stmt = stmt.filter(UmeAlarmCurrent.is_cleared == str(is_cleared).strip()) + if ne_id and str(ne_id).strip(): + stmt = stmt.filter(UmeAlarmCurrent.ne_id == str(ne_id).strip()) + kw = str(keyword or "").strip() + if kw: + 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) + ) + total = int(stmt.count()) + rows = stmt.order_by(UmeAlarmCurrent.last_seen_at.desc()).offset((page - 1) * page_size).limit(page_size).all() + 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 ""), + "event_type": str(x.event_type or ""), + "native_probable_cause": str(x.native_probable_cause or ""), + "perceived_severity": str(x.perceived_severity or ""), + "is_cleared": str(x.is_cleared or ""), + "time_created": str(x.time_created or ""), + "last_seen_at": (_ensure_utc(x.last_seen_at) or datetime.now(timezone.utc)).isoformat(), + } + for x in rows + ] + return {"total": total, "page": page, "page_size": page_size, "items": items} + + +@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) + return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne} + + +@app.get("/v1/ume/alarms/history") +def ume_list_alarms_history( + severity: str | None = Query(default=None), + ne_id: str | None = Query(default=None), + keyword: str | None = Query(default=None), + time_from: str | None = Query(default=None), + time_to: str | None = Query(default=None), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=50, ge=1, le=500), + db: Session = Depends(get_db), +) -> dict[str, Any]: + stmt = db.query(UmeAlarmHistory) + if severity and str(severity).strip(): + stmt = stmt.filter(UmeAlarmHistory.perceived_severity == str(severity).strip()) + if ne_id and str(ne_id).strip(): + stmt = stmt.filter(UmeAlarmHistory.ne_id == str(ne_id).strip()) + kw = str(keyword or "").strip() + if kw: + 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) + ) + dt_from = _parse_time(time_from) + dt_to = _parse_time(time_to) + if dt_from: + stmt = stmt.filter(UmeAlarmHistory.last_seen_at >= dt_from.replace(tzinfo=None)) + if dt_to: + stmt = stmt.filter(UmeAlarmHistory.last_seen_at <= dt_to.replace(tzinfo=None)) + total = int(stmt.count()) + rows = stmt.order_by(UmeAlarmHistory.last_seen_at.desc()).offset((page - 1) * page_size).limit(page_size).all() + 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 ""), + "event_type": str(x.event_type or ""), + "native_probable_cause": str(x.native_probable_cause or ""), + "perceived_severity": str(x.perceived_severity or ""), + "is_cleared": str(x.is_cleared or ""), + "time_created": str(x.time_created or ""), + "last_seen_at": (_ensure_utc(x.last_seen_at) or datetime.now(timezone.utc)).isoformat(), + } + for x in rows + ] + return {"total": total, "page": page, "page_size": page_size, "items": items} + + +@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]) + return {"total": len(rows), "by_severity": by_severity, "by_ne": by_ne, "by_date": by_date} + + @app.get("/v1/integrations/status") def integrations_status(db: Session = Depends(get_db)) -> dict: # netx api is up if this handler executes; still verify DB + oclaw bridge separately. diff --git a/netx_api/mcp_server.py b/netx_api/mcp_server.py index def3247..cf2d6f4 100644 --- a/netx_api/mcp_server.py +++ b/netx_api/mcp_server.py @@ -4,9 +4,28 @@ import json import sys from typing import Any -from .db import SessionLocal +from .db import Base, SessionLocal, engine from .importer import aggregate_alarms, query_alarms -from .models import AlarmBatch +from .models import AlarmBatch, UmeAlarmCurrent, UmeAlarmHistory, UmeInventoryNE +from .ume_client import UMEClient +from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full +from .ume_token_store import ( + clear_shared_token, + load_shared_token, + release_refresh_lock, + save_shared_token, + try_acquire_refresh_lock, + wait_for_token_update, +) + +_UME_CLIENT_SINGLETON = UMEClient( + token_loader=lambda: load_shared_token(), + token_saver=lambda token, exp: save_shared_token(token, exp), + token_clearer=lambda: clear_shared_token(), + lock_acquirer=lambda: try_acquire_refresh_lock(), + lock_releaser=lambda: release_refresh_lock(), + token_waiter=lambda min_exp: wait_for_token_update(min_expires_at_epoch_s=float(min_exp)), +) def _ok(rid: Any, result: dict[str, Any]) -> None: @@ -76,6 +95,75 @@ def _tool_list() -> list[dict[str, Any]]: "additionalProperties": False, }, }, + { + "name": "umeSync", + "description": "Trigger UME sync for inventory/current/history domains.", + "inputSchema": { + "type": "object", + "properties": { + "domains": { + "type": "array", + "items": {"type": "string", "enum": ["inventory", "alarms_current", "alarms_history"]}, + }, + "trigger_mode": {"type": "string", "enum": ["manual", "schedule"]}, + }, + "additionalProperties": False, + }, + }, + { + "name": "umeListNE", + "description": "List UME network elements from inventory table.", + "inputSchema": { + "type": "object", + "properties": { + "keyword": {"type": "string"}, + "page": {"type": "integer", "minimum": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500}, + }, + "additionalProperties": False, + }, + }, + { + "name": "umeGetNE", + "description": "Get UME network element by ne_id.", + "inputSchema": { + "type": "object", + "properties": {"ne_id": {"type": "string"}}, + "required": ["ne_id"], + "additionalProperties": False, + }, + }, + { + "name": "umeListCurrentAlarms", + "description": "List current UME alarms from current table.", + "inputSchema": { + "type": "object", + "properties": { + "severity": {"type": "string"}, + "is_cleared": {"type": "string"}, + "ne_id": {"type": "string"}, + "keyword": {"type": "string"}, + "page": {"type": "integer", "minimum": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500}, + }, + "additionalProperties": False, + }, + }, + { + "name": "umeListHistoryAlarms", + "description": "List historical UME alarms from history table.", + "inputSchema": { + "type": "object", + "properties": { + "severity": {"type": "string"}, + "ne_id": {"type": "string"}, + "keyword": {"type": "string"}, + "page": {"type": "integer", "minimum": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500}, + }, + "additionalProperties": False, + }, + }, ] @@ -161,10 +249,191 @@ def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: "actions": actions, } return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "umeSync": + domains_raw = args.get("domains") + domains = [] + if isinstance(domains_raw, list): + domains = [str(x).strip() for x in domains_raw if str(x).strip()] + if not domains: + domains = ["inventory", "alarms_current", "alarms_history"] + trigger_mode = str(args.get("trigger_mode") or "manual").strip().lower() + if trigger_mode not in {"manual", "schedule"}: + trigger_mode = "manual" + client = _UME_CLIENT_SINGLETON + results: list[dict[str, Any]] = [] + if "inventory" in domains: + j = sync_inventory_full(db, client, trigger_mode=trigger_mode) + results.append( + { + "domain": "inventory", + "status": j.status, + "pulled_count": int(j.pulled_count or 0), + "inserted_count": int(j.inserted_count or 0), + "updated_count": int(j.updated_count or 0), + "error_message": str(j.error_message or ""), + } + ) + if "alarms_current" in domains: + j, b = sync_alarms_current(db, client, trigger_mode=trigger_mode) + results.append( + { + "domain": "alarms_current", + "status": j.status, + "batch_id": str(b.batch_id), + "pulled_count": int(j.pulled_count or 0), + "inserted_count": int(j.inserted_count or 0), + "updated_count": int(j.updated_count or 0), + "error_message": str(j.error_message or ""), + } + ) + if "alarms_history" in domains: + j, b = sync_alarms_history_full(db, client, trigger_mode=trigger_mode) + results.append( + { + "domain": "alarms_history", + "status": j.status, + "batch_id": str(b.batch_id), + "pulled_count": int(j.pulled_count or 0), + "inserted_count": int(j.inserted_count or 0), + "updated_count": int(j.updated_count or 0), + "error_message": str(j.error_message or ""), + } + ) + return {"content": [{"type": "text", "text": json.dumps({"ok": True, "jobs": results}, ensure_ascii=False)}]} + if name == "umeListNE": + stmt = db.query(UmeInventoryNE) + keyword = str(args.get("keyword") or "").strip() + if keyword: + stmt = stmt.filter( + UmeInventoryNE.ne_id.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))) + total = int(stmt.count()) + rows = stmt.order_by(UmeInventoryNE.ne_id.asc()).offset((page - 1) * page_size).limit(page_size).all() + payload = { + "total": total, + "page": page, + "page_size": page_size, + "items": [ + { + "ne_id": str(x.ne_id or ""), + "ne_name": str(x.ne_name or ""), + "user_label": str(x.user_label or ""), + "ip_address": str(x.ip_address or ""), + "ne_type": str(x.ne_type or ""), + } + for x in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "umeGetNE": + ne_id = str(args.get("ne_id") or "").strip() + row = db.get(UmeInventoryNE, ne_id) + if not row: + return {"content": [{"type": "text", "text": json.dumps({"error": "ume_ne_not_found"})}], "isError": True} + payload = { + "ne_id": str(row.ne_id or ""), + "ne_name": str(row.ne_name or ""), + "user_label": str(row.user_label or ""), + "ip_address": str(row.ip_address or ""), + "ne_type": str(row.ne_type or ""), + "vendor": str(row.vendor or ""), + "raw_json": str(row.raw_json or "{}"), + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "umeListCurrentAlarms": + stmt = db.query(UmeAlarmCurrent) + 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() + keyword = str(args.get("keyword") or "").strip() + if severity: + stmt = stmt.filter(UmeAlarmCurrent.perceived_severity == severity) + if is_cleared: + stmt = stmt.filter(UmeAlarmCurrent.is_cleared == is_cleared) + if ne_id: + stmt = stmt.filter(UmeAlarmCurrent.ne_id == ne_id) + if keyword: + stmt = stmt.filter( + UmeAlarmCurrent.alarm_key.contains(keyword) + | UmeAlarmCurrent.object_name.contains(keyword) + | UmeAlarmCurrent.ne_name.contains(keyword) + | UmeAlarmCurrent.user_label.contains(keyword) + ) + page = max(1, int(args.get("page") or 1)) + page_size = min(500, max(1, int(args.get("page_size") or 50))) + total = int(stmt.count()) + rows = stmt.order_by(UmeAlarmCurrent.last_seen_at.desc()).offset((page - 1) * page_size).limit(page_size).all() + payload = { + "total": total, + "page": page, + "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 ""), + } + for x in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "umeListHistoryAlarms": + stmt = db.query(UmeAlarmHistory) + severity = str(args.get("severity") or "").strip() + ne_id = str(args.get("ne_id") or "").strip() + keyword = str(args.get("keyword") or "").strip() + if severity: + stmt = stmt.filter(UmeAlarmHistory.perceived_severity == severity) + if ne_id: + stmt = stmt.filter(UmeAlarmHistory.ne_id == ne_id) + if keyword: + stmt = stmt.filter( + UmeAlarmHistory.alarm_key.contains(keyword) + | UmeAlarmHistory.object_name.contains(keyword) + | UmeAlarmHistory.ne_name.contains(keyword) + | UmeAlarmHistory.user_label.contains(keyword) + ) + page = max(1, int(args.get("page") or 1)) + page_size = min(500, max(1, int(args.get("page_size") or 50))) + total = int(stmt.count()) + rows = stmt.order_by(UmeAlarmHistory.last_seen_at.desc()).offset((page - 1) * page_size).limit(page_size).all() + payload = { + "total": total, + "page": page, + "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 ""), + } + for x in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} raise ValueError(f"unknown tool: {name}") def main() -> None: + try: + Base.metadata.create_all(bind=engine) + except Exception: + pass for line in sys.stdin: raw = line.strip() if not raw: diff --git a/netx_api/models.py b/netx_api/models.py index d0417a9..ba8db9e 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -99,3 +99,111 @@ class AiAnalyzeHistory(Base): error: Mapped[str] = mapped_column(Text, default="") evidence_json: Mapped[str] = mapped_column(Text, default="{}") created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + + +class UmeSyncJob(Base): + __tablename__ = "ume_sync_jobs" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + domain: Mapped[str] = mapped_column(String(32), index=True, default="inventory") + status: Mapped[str] = mapped_column(String(32), index=True, default="running") + trigger_mode: Mapped[str] = mapped_column(String(32), default="manual") + started_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + pulled_count: Mapped[int] = mapped_column(Integer, default=0) + inserted_count: Mapped[int] = mapped_column(Integer, default=0) + updated_count: Mapped[int] = mapped_column(Integer, default=0) + error_message: Mapped[str] = mapped_column(String(1024), default="") + details_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeAlarmBatch(Base): + __tablename__ = "ume_alarm_batches" + + batch_id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + kind: Mapped[str] = mapped_column(String(16), index=True, default="current") # current/history + status: Mapped[str] = mapped_column(String(32), default="done") + total_rows: Mapped[int] = mapped_column(Integer, default=0) + success_rows: Mapped[int] = mapped_column(Integer, default=0) + failed_rows: Mapped[int] = mapped_column(Integer, default=0) + started_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + error_message: Mapped[str] = mapped_column(String(1024), default="") + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeInventoryNE(Base): + __tablename__ = "ume_inventory_ne" + + ne_id: Mapped[str] = mapped_column(String(128), primary_key=True) + ne_name: Mapped[str] = mapped_column(String(256), default="", index=True) + user_label: Mapped[str] = mapped_column(String(256), default="", index=True) + ip_address: Mapped[str] = mapped_column(String(128), default="", index=True) + ne_type: Mapped[str] = mapped_column(String(128), default="") + vendor: Mapped[str] = mapped_column(String(64), default="ZTE") + source_type: Mapped[str] = mapped_column(String(64), default="ume_restconf") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeInventoryEquipmentHolder(Base): + __tablename__ = "ume_inventory_equipment_holder" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + ne_id: Mapped[str] = mapped_column(ForeignKey("ume_inventory_ne.ne_id"), index=True) + holder_name: Mapped[str] = mapped_column(String(256), index=True, default="") + holder_type: Mapped[str] = mapped_column(String(128), default="") + holder_state: Mapped[str] = mapped_column(String(128), default="") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeAlarmCurrent(Base): + __tablename__ = "ume_alarms_current" + + 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="") + perceived_severity: Mapped[str] = mapped_column(String(64), default="", index=True) + is_cleared: Mapped[str] = mapped_column(String(16), default="", index=True) + time_created: Mapped[str] = mapped_column(String(64), default="", index=True) + root_cause_alarm_indication: Mapped[str] = mapped_column(String(32), default="") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeAlarmHistory(Base): + __tablename__ = "ume_alarms_history" + + 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="") + perceived_severity: Mapped[str] = mapped_column(String(64), default="", index=True) + is_cleared: Mapped[str] = mapped_column(String(16), default="", index=True) + time_created: Mapped[str] = mapped_column(String(64), default="", index=True) + root_cause_alarm_indication: Mapped[str] = mapped_column(String(32), default="") + first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + raw_json: Mapped[str] = mapped_column(Text, default="{}") + + +class UmeTokenCache(Base): + __tablename__ = "ume_token_cache" + + cache_key: Mapped[str] = mapped_column(String(256), primary_key=True) + token: Mapped[str] = mapped_column(Text, default="") + expires_at_epoch_s: Mapped[int] = mapped_column(Integer, default=0) + lock_owner: Mapped[str] = mapped_column(String(128), default="", index=True) + lock_expires_at_epoch_s: Mapped[int] = mapped_column(Integer, default=0, index=True) + updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) diff --git a/netx_api/ume_client.py b/netx_api/ume_client.py new file mode 100644 index 0000000..9b281ad --- /dev/null +++ b/netx_api/ume_client.py @@ -0,0 +1,458 @@ +from __future__ import annotations + +from dataclasses import dataclass +import json +from threading import Lock +from time import time +from typing import Any, Callable + +import httpx + +from .config import settings + + +def _rstrip_slash(url: str) -> str: + return str(url or "").strip().rstrip("/") + + +def _coerce_dict(value: Any) -> dict[str, Any]: + return value if isinstance(value, dict) else {} + + +def _coerce_list(value: Any) -> list[Any]: + return value if isinstance(value, list) else [] + + +@dataclass +class RequestDiagnostics: + method: str + path: str + status_code: int + latency_ms: int + retry_count: int = 0 + error_code: str = "" + + +class UMEClient: + def __init__( + self, + *, + base_url: str | None = None, + username: str | None = None, + password: str | None = None, + verify_tls: bool | None = None, + timeout_s: float | None = None, + auth_header: str | None = None, + content_type: str | None = None, + token_ttl_s: int | None = None, + token_refresh_skew_s: int | None = None, + token_path: str | None = None, + token_handshake_path: str | None = None, + token_logout_path: str | None = None, + ne_path: str | None = None, + alarms_path: str | None = None, + token_loader: Callable[[], tuple[str, float] | None] | None = None, + token_saver: Callable[[str, float], None] | None = None, + token_clearer: Callable[[], None] | None = None, + lock_acquirer: Callable[[], bool] | None = None, + lock_releaser: Callable[[], None] | None = None, + token_waiter: Callable[[float], tuple[str, float] | None] | None = None, + ) -> None: + self.base_url = _rstrip_slash(base_url if base_url is not None else settings.ume_base_url) + self.username = str(username if username is not None else settings.ume_username) + self.password = str(password if password is not None else settings.ume_password) + self.verify_tls = bool(settings.ume_verify_tls if verify_tls is None else verify_tls) + self.timeout_s = max(3.0, float(settings.ume_timeout_s if timeout_s is None else timeout_s)) + self.auth_header = str(auth_header if auth_header is not None else settings.ume_auth_header).strip() or "accessToken" + self.content_type = str(content_type if content_type is not None else settings.ume_content_type).strip() + self.token_ttl_s = max(60, int(settings.ume_token_ttl_s if token_ttl_s is None else token_ttl_s)) + self.token_refresh_skew_s = max( + 5, int(settings.ume_token_refresh_skew_s if token_refresh_skew_s is None else token_refresh_skew_s) + ) + self.token_path = str(token_path if token_path is not None else settings.ume_token_path).strip() + self.token_handshake_path = str( + token_handshake_path if token_handshake_path is not None else settings.ume_token_handshake_path + ).strip() + self.token_logout_path = str(token_logout_path if token_logout_path is not None else settings.ume_token_logout_path).strip() + self.ne_path = str(ne_path if ne_path is not None else settings.ume_ne_path).strip() + self.alarms_path = str(alarms_path if alarms_path is not None else settings.ume_alarms_path).strip() + self._token_loader = token_loader + self._token_saver = token_saver + self._token_clearer = token_clearer + self._lock_acquirer = lock_acquirer + self._lock_releaser = lock_releaser + self._token_waiter = token_waiter + + self._token_lock = Lock() + self._token_value: str = "" + self._token_expires_at: float = 0.0 + self._last_token_source: str = "memory" + self._last_store_sync_changed: bool = False + + def _sync_token_from_store(self) -> None: + self._last_store_sync_changed = False + if self._token_loader is None: + return + try: + loaded = self._token_loader() + except Exception: + return + if not loaded: + return + token, exp = loaded + token = str(token or "").strip() + if not token: + return + if exp > float(self._token_expires_at): + self._token_value = token + self._token_expires_at = float(exp) + self._last_token_source = "db" + self._last_store_sync_changed = True + + def _persist_token_to_store(self) -> None: + if self._token_saver is None: + return + try: + self._token_saver(self._token_value, self._token_expires_at) + except Exception: + pass + + def _clear_token_in_store(self) -> None: + if self._token_clearer is None: + return + try: + self._token_clearer() + except Exception: + pass + + def token_status(self) -> dict[str, Any]: + self._sync_token_from_store() + now = time() + token = self._token_value.strip() + has_token = bool(token) + expires_in_s = int(max(0, self._token_expires_at - now)) if has_token else 0 + token_preview = "" + if token: + token_preview = f"{token[:8]}...{token[-4:]}" if len(token) > 16 else token + return { + "has_token": has_token, + "expires_in_s": expires_in_s, + "expires_at_epoch_s": int(self._token_expires_at) if has_token else 0, + "auth_header": self.auth_header, + "token_preview": token_preview, + "source": str(self._last_token_source or "memory"), + "store_synced": bool(self._last_store_sync_changed), + } + + def _assert_ready(self) -> None: + if not self.base_url: + raise RuntimeError("ume_base_url_required") + if not self.username: + raise RuntimeError("ume_username_required") + if not self.password: + raise RuntimeError("ume_password_required") + + def _build_url(self, path: str) -> str: + p = str(path or "").strip() + if not p: + raise RuntimeError("ume_path_required") + if p.startswith("http://") or p.startswith("https://"): + return p + if not p.startswith("/"): + p = "/" + p + return f"{self.base_url}{p}" + + def _headers(self, *, include_token: bool = True) -> dict[str, str]: + headers = { + "accept": "application/yang-data+json", + "content-type": self.content_type, + } + if include_token: + token = self._token_value.strip() + if token: + headers[self.auth_header] = token + return headers + + def _extract_token_and_ttl(self, payload: dict[str, Any]) -> tuple[str, int | None]: + token = "" + ttl: int | None = None + + def walk(node: Any) -> None: + nonlocal token, ttl + if isinstance(node, dict): + for k, v in node.items(): + key = str(k).lower() + if not token and key in {"accesstoken", "access_token", "token"} and isinstance(v, str): + token = v.strip() + if ttl is None and key in {"expires", "expiresin", "expires_in", "ttl"}: + try: + ttl = int(v) + except Exception: + pass + walk(v) + elif isinstance(node, list): + for item in node: + walk(item) + + walk(payload) + return token, ttl + + def login(self, *, force: bool = False) -> str: + self._assert_ready() + self._sync_token_from_store() + now = time() + if not force and self._token_value and now < (self._token_expires_at - self.token_refresh_skew_s): + return self._token_value + + with self._token_lock: + now = time() + if not force and self._token_value and now < (self._token_expires_at - self.token_refresh_skew_s): + return self._token_value + + # Cross-process singleflight: if another process is refreshing, wait for DB update. + min_exp = float(self._token_expires_at) + if self._lock_acquirer is not None: + acquired = False + try: + acquired = bool(self._lock_acquirer()) + except Exception: + acquired = False + if not acquired: + if self._token_waiter is not None: + waited = self._token_waiter(min_exp) + if waited: + self._sync_token_from_store() + now = time() + if self._token_value and now < (self._token_expires_at - self.token_refresh_skew_s): + return self._token_value + + url = self._build_url(self.token_path) + payload = {"login-info": {"user-name": self.username, "value": self.password}} + try: + with httpx.Client(verify=self.verify_tls, timeout=self.timeout_s) as client: + resp = client.post(url, json=payload, headers=self._headers(include_token=False)) + if not resp.is_success: + raise RuntimeError(f"ume_login_failed:{resp.status_code}:{resp.text[:240]}") + data = _coerce_dict(resp.json()) + except Exception as exc: + if self._lock_releaser is not None: + try: + self._lock_releaser() + except Exception: + pass + raise RuntimeError(f"ume_login_failed:{str(exc)[:240]}") from exc + + token, ttl = self._extract_token_and_ttl(data) + if not token: + if self._lock_releaser is not None: + try: + self._lock_releaser() + except Exception: + pass + raise RuntimeError("ume_login_failed:missing_access_token") + use_ttl = max(60, int(ttl)) if ttl is not None else self.token_ttl_s + self._token_value = token + self._token_expires_at = time() + use_ttl + self._last_token_source = "memory" + self._persist_token_to_store() + if self._lock_releaser is not None: + try: + self._lock_releaser() + except Exception: + pass + return self._token_value + + def renew_token(self) -> str: + self._assert_ready() + token = self._token_value.strip() + if not token: + return self.login(force=True) + with self._token_lock: + token = self._token_value.strip() + if not token: + return self.login(force=True) + url = self._build_url(self.token_handshake_path) + try: + with httpx.Client(verify=self.verify_tls, timeout=self.timeout_s) as client: + resp = client.post(url, headers=self._headers(include_token=True)) + if not resp.is_success: + # handshake may fail if token expired; fallback to full login + return self.login(force=True) + # Per UME guide, oauth_handshake may return no body; treat HTTP 2xx as success. + next_token = "" + ttl: int | None = None + try: + text = (resp.text or "").strip() + if text: + data = _coerce_dict(resp.json()) + next_token, ttl = self._extract_token_and_ttl(data) + except Exception: + # ignore json parse errors, success is based on status code + next_token = "" + ttl = None + if next_token: + self._token_value = next_token + # If handshake doesn't provide ttl, fall back to configured ttl. + use_ttl = max(60, int(ttl)) if ttl is not None else self.token_ttl_s + self._token_expires_at = time() + use_ttl + self._last_token_source = "memory" + self._persist_token_to_store() + return self._token_value + except Exception: + return self.login(force=True) + + def logout_token(self) -> bool: + token = self._token_value.strip() + if not token: + return True + url = self._build_url(self.token_logout_path) + try: + with httpx.Client(verify=self.verify_tls, timeout=self.timeout_s) as client: + resp = client.delete(url, headers=self._headers(include_token=True)) + if resp.status_code in (401, 403): + # Token already invalid/expired on server side can be treated as logged out. + self._token_value = "" + self._token_expires_at = 0.0 + self._last_token_source = "memory" + self._clear_token_in_store() + return True + if not resp.is_success: + return False + self._token_value = "" + self._token_expires_at = 0.0 + self._last_token_source = "memory" + self._clear_token_in_store() + return True + except Exception: + return False + + def refresh_if_needed(self) -> str: + self._sync_token_from_store() + now = time() + if self._token_value and now < (self._token_expires_at - self.token_refresh_skew_s): + return self._token_value + if self._token_value: + return self.renew_token() + return self.login(force=False) + + def request_json( + self, + method: str, + path: str, + *, + params: dict[str, Any] | None = None, + body: dict[str, Any] | None = None, + ) -> tuple[dict[str, Any], RequestDiagnostics]: + self.refresh_if_needed() + url = self._build_url(path) + m = str(method or "GET").upper() + retry_count = 0 + t0 = time() + try: + with httpx.Client(verify=self.verify_tls, timeout=self.timeout_s) as client: + resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True)) + if resp.status_code in (401, 403): + retry_count = 1 + self.login(force=True) + with httpx.Client(verify=self.verify_tls, timeout=self.timeout_s) as client: + resp = client.request(m, url, params=params, json=body, headers=self._headers(include_token=True)) + if not resp.is_success: + diag = RequestDiagnostics( + method=m, + path=path, + status_code=int(resp.status_code), + latency_ms=int((time() - t0) * 1000), + retry_count=retry_count, + error_code=f"http_{int(resp.status_code)}", + ) + raise RuntimeError(f"ume_request_failed:{resp.status_code}:{resp.text[:240]}") + data = _coerce_dict(resp.json()) + diag = RequestDiagnostics( + method=m, + path=path, + status_code=int(resp.status_code), + latency_ms=int((time() - t0) * 1000), + retry_count=retry_count, + ) + return data, diag + except Exception as exc: + if isinstance(exc, RuntimeError): + raise + raise RuntimeError(f"ume_request_failed:{str(exc)[:240]}") from exc + + def _extract_named_list(self, payload: dict[str, Any], keys: list[str]) -> list[dict[str, Any]]: + target_keys = {str(k).lower() for k in keys} + found: list[dict[str, Any]] = [] + seen: set[str] = set() + + def add_unique(item: dict[str, Any]) -> None: + try: + mark = json.dumps(item, ensure_ascii=False, sort_keys=True, default=str) + except Exception: + mark = str(item) + if mark in seen: + return + seen.add(mark) + found.append(item) + + def walk(node: Any) -> None: + if isinstance(node, dict): + for k, v in node.items(): + if str(k).lower() in target_keys: + if isinstance(v, list): + for it in v: + if isinstance(it, dict): + add_unique(it) + elif isinstance(v, dict): + # Common RESTCONF wrappers are container dicts, e.g. + # network-elements -> network-element[] / alarm-list -> alarm[]. + # Prefer unwrapping nested list payloads before falling back. + nested_collected = False + for nested_v in v.values(): + if isinstance(nested_v, list): + for it in nested_v: + if isinstance(it, dict): + add_unique(it) + nested_collected = True + if not nested_collected: + add_unique(v) + walk(v) + elif isinstance(node, list): + for item in node: + walk(item) + + walk(payload) + return found + + def get_network_elements(self) -> tuple[list[dict[str, Any]], RequestDiagnostics]: + data, diag = self.request_json("GET", self.ne_path) + rows = self._extract_named_list(data, ["network-elements", "network-element", "ne", "network_elements"]) + if rows: + return rows, diag + # fallback: some responses may directly return list-like map at top-level + for v in data.values(): + lst = _coerce_list(v) + if lst and isinstance(lst[0], dict): + return [x for x in lst if isinstance(x, dict)], diag + return [], diag + + def get_alarms( + self, + *, + is_uncleared: bool, + limit: int | None = None, + offset: int | 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] = { + "is-uncleared": "true" if is_uncleared else "false", + "limit": page_size, + } + if offset is not None and int(offset) >= 0: + params["offset"] = int(offset) + 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 new file mode 100644 index 0000000..c1e0e2c --- /dev/null +++ b/netx_api/ume_sync_service.py @@ -0,0 +1,299 @@ +from __future__ import annotations + +import json +from datetime import datetime, timedelta, timezone +from typing import Any + +from sqlalchemy.orm import Session + +from .config import settings +from .models import ( + UmeAlarmBatch, + UmeAlarmCurrent, + UmeAlarmHistory, + UmeInventoryEquipmentHolder, + UmeInventoryNE, + UmeSyncJob, +) +from .ume_client import UMEClient + + +def _s(v: Any) -> str: + if v is None: + return "" + text = str(v).strip() + if text.lower() == "nan": + return "" + return text + + +def _utc_now_naive() -> datetime: + return datetime.now(timezone.utc).replace(tzinfo=None) + + +def _pick(d: dict[str, Any], *keys: str) -> Any: + for key in keys: + if key in d: + return d.get(key) + return None + + +def _alarm_key(alarm: dict[str, Any]) -> str: + key = _s(_pick(alarm, "alarmKey", "alarm-key", "id")) + if key: + return key + parts = [ + _s(_pick(alarm, "objectName", "object-name")), + _s(_pick(alarm, "eventType", "event-type")), + _s(_pick(alarm, "timeCreated", "time-created")), + _s(_pick(alarm, "nativeProbableCause", "native-probable-cause")), + ] + merged = "|".join(x for x in parts if x) + return merged or f"fallback-{datetime.utcnow().timestamp()}" + + +def _build_sync_job(domain: str, trigger_mode: str) -> UmeSyncJob: + return UmeSyncJob( + domain=domain, + status="running", + trigger_mode=trigger_mode, + started_at=_utc_now_naive(), + ) + + +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() + now = _utc_now_naive() + pulled = len(ne_rows) + for row in ne_rows: + ne_id = _s(_pick(row, "ne-id", "ne_id", "id")) + if not ne_id: + continue + existing = db.get(UmeInventoryNE, ne_id) + if existing is None: + existing = UmeInventoryNE( + ne_id=ne_id, + first_seen_at=now, + ) + db.add(existing) + inserted += 1 + else: + updated += 1 + existing.ne_name = _s(_pick(row, "name", "ne-name")) + existing.user_label = _s(_pick(row, "user-label", "user_label")) + existing.ip_address = _s(_pick(row, "ip-Address", "ip-address", "ip")) + existing.ne_type = _s(_pick(row, "type", "ne-type")) + existing.last_seen_at = now + existing.raw_json = json.dumps(row, ensure_ascii=False, default=str) + + # Optional: holders may not exist in current UME deployment; keep best-effort. + # This pass keeps schema warm for later detailed holder endpoint integration. + holder_payload: list[dict[str, Any]] = [] + for ne in ne_rows: + holders = _pick(ne, "equipment-holder", "equipment-holders") + if isinstance(holders, list): + for h in holders: + if isinstance(h, dict): + h2 = dict(h) + if "ne-id" not in h2: + h2["ne-id"] = _pick(ne, "ne-id", "ne_id", "id") + holder_payload.append(h2) + for h in holder_payload: + ne_id = _s(_pick(h, "ne-id", "ne_id")) + holder_name = _s(_pick(h, "name", "holder-name")) + if not ne_id or not holder_name: + continue + existing_holder = ( + db.query(UmeInventoryEquipmentHolder) + .filter( + UmeInventoryEquipmentHolder.ne_id == ne_id, + UmeInventoryEquipmentHolder.holder_name == holder_name, + ) + .one_or_none() + ) + if existing_holder is None: + existing_holder = UmeInventoryEquipmentHolder( + ne_id=ne_id, + holder_name=holder_name, + first_seen_at=now, + ) + db.add(existing_holder) + existing_holder.holder_type = _s(_pick(h, "type", "holder-type")) + existing_holder.holder_state = _s(_pick(h, "state", "holder-state")) + existing_holder.last_seen_at = now + existing_holder.raw_json = json.dumps(h, ensure_ascii=False, default=str) + + job.status = "done" + except Exception as exc: + job.status = "failed" + job.error_message = str(exc)[:1024] + finally: + job.pulled_count = int(pulled) + job.inserted_count = int(inserted) + job.updated_count = int(updated) + job.ended_at = _utc_now_naive() + db.commit() + db.refresh(job) + return job + + +def _sync_alarms_common( + db: Session, + client: UMEClient, + *, + is_uncleared: bool, + trigger_mode: str, +) -> tuple[UmeSyncJob, UmeAlarmBatch]: + domain = "alarms_history" if is_uncleared else "alarms_current" + job = _build_sync_job(domain, trigger_mode) + db.add(job) + batch = UmeAlarmBatch( + kind="history" if is_uncleared else "current", + status="running", + started_at=_utc_now_naive(), + ) + db.add(batch) + db.flush() + now = _utc_now_naive() + pulled = inserted = updated = 0 + paging_mode = "offset" + 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 = max(1, min(page_size, limit_max)) + max_pages = int(getattr(settings, "ume_max_pages", 2000) or 2000) + max_pages = max(1, min(max_pages, 20000)) + offset = 0 + page_no = 0 + + def upsert_alarm(alarm: dict[str, Any]) -> 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) + db.add(existing) + inserted += 1 + else: + updated += 1 + else: + existing = db.get(UmeAlarmCurrent, key) + if existing is None: + existing = UmeAlarmCurrent(alarm_key=key, first_seen_at=now) + db.add(existing) + inserted += 1 + else: + updated += 1 + existing.ne_id = _s(_pick(alarm, "ne-id", "neId", "ne_id")) + 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")) + existing.perceived_severity = _s(_pick(alarm, "perceivedSeverity", "perceived-severity")) + existing.is_cleared = _s(_pick(alarm, "isCleared", "is-cleared")) + existing.time_created = _s(_pick(alarm, "timeCreated", "time-created")) + existing.root_cause_alarm_indication = _s( + _pick(alarm, "rootCauseAlarmIndication", "root-cause-alarm-indication") + ) + 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) + + # 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: + 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 + + if not is_uncleared: + expiry = now - timedelta(hours=48) + ( + db.query(UmeAlarmCurrent) + .filter(UmeAlarmCurrent.last_seen_at < expiry) + .delete(synchronize_session=False) + ) + + batch.total_rows = int(pulled) + batch.success_rows = int(inserted + updated) + batch.failed_rows = max(0, int(pulled) - int(inserted + updated)) + batch.status = "done" + batch.ended_at = _utc_now_naive() + batch.raw_json = json.dumps( + {"pulled": pulled, "inserted": inserted, "updated": updated, "paging_mode": paging_mode}, + ensure_ascii=False, + ) + + job.status = "done" + except Exception as exc: + msg = str(exc)[:1024] + batch.status = "failed" + batch.error_message = msg + batch.ended_at = _utc_now_naive() + job.status = "failed" + job.error_message = msg + finally: + job.pulled_count = int(pulled) + job.inserted_count = int(inserted) + job.updated_count = int(updated) + job.ended_at = _utc_now_naive() + job.details_json = json.dumps( + { + "batch_id": batch.batch_id, + "kind": batch.kind, + "status": batch.status, + "paging_mode": paging_mode, + "paging_note": paging_note, + }, + ensure_ascii=False, + ) + db.commit() + db.refresh(job) + db.refresh(batch) + return job, batch + + +def sync_alarms_current(db: Session, client: UMEClient, *, trigger_mode: str = "manual") -> tuple[UmeSyncJob, UmeAlarmBatch]: + return _sync_alarms_common(db, client, is_uncleared=False, trigger_mode=trigger_mode) + + +def sync_alarms_history_full( + db: Session, client: UMEClient, *, trigger_mode: str = "manual" +) -> tuple[UmeSyncJob, UmeAlarmBatch]: + return _sync_alarms_common(db, client, is_uncleared=True, trigger_mode=trigger_mode) diff --git a/netx_api/ume_token_store.py b/netx_api/ume_token_store.py new file mode 100644 index 0000000..966cac3 --- /dev/null +++ b/netx_api/ume_token_store.py @@ -0,0 +1,143 @@ +from __future__ import annotations + +import os +import socket +from time import sleep, time +from datetime import datetime + +from sqlalchemy import text as sql_text + +from .config import settings +from .db import SessionLocal +from .models import UmeTokenCache + + +def build_cache_key() -> str: + base = str(settings.ume_base_url or "").strip().rstrip("/") + user = str(settings.ume_username or "").strip() + return f"{base}|{user}" if base or user else "default" + +def build_owner_id() -> str: + host = "" + try: + host = socket.gethostname() + except Exception: + host = "host" + pid = os.getpid() + return f"{host}:{pid}" + + +def load_shared_token(cache_key: str | None = None) -> tuple[str, float] | None: + key = str(cache_key or build_cache_key()) + with SessionLocal() as db: + row = db.get(UmeTokenCache, key) + if not row: + return None + token = str(row.token or "").strip() + if not token: + return None + return token, float(int(row.expires_at_epoch_s or 0)) + +def try_acquire_refresh_lock( + *, + cache_key: str | None = None, + owner_id: str | None = None, + lock_ttl_s: int = 30, +) -> bool: + key = str(cache_key or build_cache_key()) + owner = str(owner_id or build_owner_id()) + now = int(time()) + ttl = max(5, int(lock_ttl_s or 30)) + lock_until = now + ttl + with SessionLocal() as db: + db.execute( + sql_text( + """ +UPDATE ume_token_cache +SET lock_owner = :owner, + lock_expires_at_epoch_s = :lock_until, + updated_at = updated_at +WHERE cache_key = :key + AND ( + lock_expires_at_epoch_s IS NULL + OR lock_expires_at_epoch_s < :now + OR lock_owner = :owner + ) +""" + ), + {"owner": owner, "lock_until": lock_until, "key": key, "now": now}, + ) + if db.get(UmeTokenCache, key) is None: + row = UmeTokenCache(cache_key=key) + db.add(row) + db.commit() + row = db.get(UmeTokenCache, key) + if not row: + return False + ok = (str(row.lock_owner or "") == owner) and int(row.lock_expires_at_epoch_s or 0) >= now + if ok: + db.commit() + return bool(ok) + + +def release_refresh_lock(*, cache_key: str | None = None, owner_id: str | None = None) -> None: + key = str(cache_key or build_cache_key()) + owner = str(owner_id or build_owner_id()) + with SessionLocal() as db: + db.execute( + sql_text( + """ +UPDATE ume_token_cache +SET lock_owner = '', + lock_expires_at_epoch_s = 0 +WHERE cache_key = :key AND lock_owner = :owner +""" + ), + {"key": key, "owner": owner}, + ) + db.commit() + + +def wait_for_token_update( + *, + cache_key: str | None = None, + min_expires_at_epoch_s: float = 0.0, + timeout_s: float = 10.0, + poll_interval_s: float = 0.4, +) -> tuple[str, float] | None: + key = str(cache_key or build_cache_key()) + deadline = time() + max(0.1, float(timeout_s or 10.0)) + poll = max(0.1, float(poll_interval_s or 0.4)) + while time() < deadline: + cur = load_shared_token(key) + if cur: + token, exp = cur + if float(exp) > float(min_expires_at_epoch_s) and str(token or "").strip(): + return token, float(exp) + sleep(poll) + return None + + +def save_shared_token(token: str, expires_at_epoch_s: float, cache_key: str | None = None) -> None: + key = str(cache_key or build_cache_key()) + with SessionLocal() as db: + row = db.get(UmeTokenCache, key) + if row is None: + row = UmeTokenCache(cache_key=key) + db.add(row) + row.token = str(token or "") + row.expires_at_epoch_s = int(expires_at_epoch_s or 0) + row.updated_at = datetime.utcnow() + db.commit() + + +def clear_shared_token(cache_key: str | None = None) -> None: + key = str(cache_key or build_cache_key()) + with SessionLocal() as db: + row = db.get(UmeTokenCache, key) + if row is None: + return + row.token = "" + row.expires_at_epoch_s = 0 + row.updated_at = datetime.utcnow() + db.commit() diff --git a/scripts/stop_netx.ps1 b/scripts/stop_netx.ps1 index 72eddcd..edf7699 100644 --- a/scripts/stop_netx.ps1 +++ b/scripts/stop_netx.ps1 @@ -31,8 +31,9 @@ function Get-ListenPids { Write-Host "[WARN] port query failed for $LocalPort : $($_.Exception.Message)" } finally { if ($job) { - Stop-Job $job -Force -ErrorAction SilentlyContinue - Remove-Job $job -Force -ErrorAction SilentlyContinue + # Windows PowerShell 5.1 does not support -Force on Stop-Job/Remove-Job. + Stop-Job $job -ErrorAction SilentlyContinue + Remove-Job $job -ErrorAction SilentlyContinue } } @($ids) diff --git a/tests/test_ume_sync.py b/tests/test_ume_sync.py new file mode 100644 index 0000000..18c41c8 --- /dev/null +++ b/tests/test_ume_sync.py @@ -0,0 +1,255 @@ +from __future__ import annotations + +import unittest +from unittest.mock import patch + +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker + +from netx_api.db import Base +from netx_api.models import UmeAlarmCurrent, UmeInventoryNE +from netx_api.ume_client import UMEClient +from netx_api.ume_sync_service import sync_alarms_current, sync_inventory_full + + +class _FakeResponse: + def __init__(self, status_code: int, payload: dict): + self.status_code = int(status_code) + self._payload = payload + self.text = str(payload) + self.is_success = 200 <= self.status_code < 300 + + def json(self): + return self._payload + + +class _FakeClient: + def __init__(self, responses): + self._responses = list(responses) + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc, tb): + return False + + def post(self, url, json=None, headers=None): + if not self._responses: + raise RuntimeError("no_more_fake_responses") + return self._responses.pop(0) + + def request(self, method, url, params=None, json=None, headers=None): + if not self._responses: + raise RuntimeError("no_more_fake_responses") + return self._responses.pop(0) + + +class UMEClientTests(unittest.TestCase): + def test_login_and_cached_token(self): + responses = [ + _FakeResponse(200, {"output": {"accessToken": "token-1", "expires": 1800}}), + ] + with patch("netx_api.ume_client.httpx.Client", return_value=_FakeClient(responses)): + client = UMEClient( + base_url="https://ume.local:18014", + username="u", + password="p", + verify_tls=False, + ) + t1 = client.login() + t2 = client.login() + self.assertEqual(t1, "token-1") + self.assertEqual(t2, "token-1") + + def test_request_retry_after_401(self): + responses = [ + _FakeResponse(200, {"output": {"accessToken": "token-1", "expires": 1800}}), + _FakeResponse(401, {"error": "expired"}), + _FakeResponse(200, {"output": {"accessToken": "token-2", "expires": 1800}}), + _FakeResponse(200, {"alarm-list": []}), + ] + with patch("netx_api.ume_client.httpx.Client", return_value=_FakeClient(responses)): + client = UMEClient( + base_url="https://ume.local:18014", + username="u", + password="p", + verify_tls=False, + ) + payload, diag = client.request_json("GET", "/restconf/data/zte-alarms:alarms/alarm-list") + self.assertIsInstance(payload, dict) + self.assertEqual(diag.retry_count, 1) + self.assertEqual(client.refresh_if_needed(), "token-2") + + def test_extract_network_elements_from_wrapped_container(self): + client = UMEClient(base_url="https://ume.local:18014", username="u", password="p", verify_tls=False) + + def _fake_request_json(method: str, path: str, **kwargs): + return ( + { + "zte-resources-module:network-elements": { + "network-element": [ + {"ne-id": "NE-1", "name": "ne1"}, + {"ne-id": "NE-2", "name": "ne2"}, + ] + } + }, + None, + ) + + client.request_json = _fake_request_json # type: ignore[method-assign] + rows, _ = client.get_network_elements() + self.assertEqual(len(rows), 2) + self.assertEqual(rows[0].get("ne-id"), "NE-1") + + def test_extract_alarm_rows_from_wrapped_alarm_list(self): + client = UMEClient(base_url="https://ume.local:18014", username="u", password="p", verify_tls=False) + + def _fake_request_json(method: str, path: str, **kwargs): + return ( + { + "zte-alarms:alarms": { + "alarm-list": { + "alarm": [ + {"alarmKey": "AK-1", "ne-id": "NE-1"}, + {"alarmKey": "AK-2", "ne-id": "NE-2"}, + ] + } + } + }, + None, + ) + + client.request_json = _fake_request_json # type: ignore[method-assign] + rows, _ = client.get_alarms(is_uncleared=False) + self.assertEqual(len(rows), 2) + self.assertEqual(rows[1].get("alarmKey"), "AK-2") + + +class UmeSyncServiceTests(unittest.TestCase): + def setUp(self): + engine = create_engine("sqlite+pysqlite:///:memory:", future=True) + TestingSessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False, expire_on_commit=False) + Base.metadata.create_all(bind=engine) + self.db = TestingSessionLocal() + + def tearDown(self): + self.db.close() + + def test_sync_inventory_upsert(self): + class _C: + def get_network_elements(self): + rows = [ + {"ne-id": "NE-1", "name": "ne1", "user-label": "网元1", "ip-Address": "10.0.0.1", "type": "A"}, + ] + return rows, None + + job1 = sync_inventory_full(self.db, _C(), trigger_mode="manual") + self.assertEqual(job1.status, "done") + self.assertEqual(job1.inserted_count, 1) + + job2 = sync_inventory_full(self.db, _C(), trigger_mode="manual") + self.assertEqual(job2.status, "done") + self.assertEqual(job2.updated_count, 1) + + ne = self.db.get(UmeInventoryNE, "NE-1") + self.assertIsNotNone(ne) + self.assertEqual(ne.user_label, "网元1") + + def test_sync_current_alarms_upsert(self): + class _C: + def get_alarms(self, *, is_uncleared: bool, limit=None, offset=None): + rows = [ + { + "alarmKey": "AK-1", + "ne-id": "NE-1", + "ne-name": "ne1", + "user-label": "网元1", + "objectName": "port-1", + "eventType": "COMMUNICATION", + "nativeProbableCause": "LOS", + "perceivedSeverity": "critical", + "isCleared": "false", + "timeCreated": "2026-01-01T00:00:00Z", + } + ] + return rows, None + + job1, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + self.assertEqual(job1.status, "done") + self.assertEqual(job1.inserted_count, 1) + + job2, _ = sync_alarms_current(self.db, _C(), trigger_mode="manual") + self.assertEqual(job2.status, "done") + self.assertEqual(job2.updated_count, 1) + + row = self.db.get(UmeAlarmCurrent, "AK-1") + self.assertIsNotNone(row) + self.assertEqual(row.perceived_severity, "critical") + + def test_sync_current_alarms_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: + return ( + [ + {"alarmKey": "AK-1", "ne-id": "NE-1", "perceivedSeverity": "major", "isCleared": "false"}, + {"alarmKey": "AK-2", "ne-id": "NE-2", "perceivedSeverity": "major", "isCleared": "false"}, + ], + None, + ) + if off == 2: + return ( + [ + {"alarmKey": "AK-3", "ne-id": "NE-3", "perceivedSeverity": "minor", "isCleared": "false"}, + ], + None, + ) + return ([], None) + + 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 = 2 + svc.settings.ume_max_pages = 10 + 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 + self.assertEqual(job.status, "done") + self.assertEqual(job.pulled_count, 3) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-3")) + + def test_sync_current_alarms_offset_unsupported_fallback(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, + ) + + 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 + 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 + + self.assertEqual(job.status, "done") + self.assertEqual(job.pulled_count, 1) + self.assertIsNotNone(self.db.get(UmeAlarmCurrent, "AK-1")) + + +if __name__ == "__main__": + unittest.main() diff --git a/web/src/App.tsx b/web/src/App.tsx index f70da06..4d037ee 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -5,8 +5,8 @@ import { AppLayout } from "./layout/AppLayout"; import { AlarmsPage } from "./pages/AlarmsPage"; import { DiagnosticsPage } from "./pages/DiagnosticsPage"; import { AiAnalysisPage } from "./pages/AiAnalysisPage"; -import { JobsPage } from "./pages/JobsPage"; import { IngestPage } from "./pages/IngestPage"; +import { UmePage } from "./pages/UmePage"; import type { Alarm, Batch, Diagnostics, ImportHistoryItem } from "./types"; import { apiDelete, apiPost, apiPostForm, fetchAlarms, fetchBatches, fetchDiagnostics, fetchIntegrationStatus } from "./services/api"; import { fetchJobs } from "./services/api"; @@ -17,7 +17,6 @@ function App() { alarms: (import.meta.env.VITE_FEATURE_ALARMS ?? "1") !== "0", diagnostics: (import.meta.env.VITE_FEATURE_DIAGNOSTICS ?? "1") !== "0", aiAnalysis: (import.meta.env.VITE_FEATURE_AI_ANALYSIS ?? "1") !== "0", - jobs: (import.meta.env.VITE_FEATURE_JOBS ?? "0") === "1", }; const [searchParams, setSearchParams] = useSearchParams(); @@ -336,7 +335,7 @@ function App() { /> } /> - } /> + } /> {toast &&
{toast.text}
} diff --git a/web/src/layout/AppLayout.tsx b/web/src/layout/AppLayout.tsx index b7445ae..f68e50a 100644 --- a/web/src/layout/AppLayout.tsx +++ b/web/src/layout/AppLayout.tsx @@ -16,7 +16,6 @@ type Props = { alarms: boolean; diagnostics: boolean; aiAnalysis: boolean; - jobs: boolean; }; children: ReactNode; }; @@ -52,11 +51,8 @@ export function AppLayout({ status, connections, onRefreshBatches, onRefreshAlar > AI 分析 - `nav-item${isActive ? " active" : ""}${!capabilities.jobs ? " disabled" : ""}`} - > - 作业中心 + `nav-item${isActive ? " active" : ""}`}> + UME 对接 diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx new file mode 100644 index 0000000..1878542 --- /dev/null +++ b/web/src/pages/UmePage.tsx @@ -0,0 +1,486 @@ +import { useState } from "react"; +import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; +import { + apiPost, + disconnectUmeToken, + fetchUmeCurrentAlarms, + fetchUmeHistoryAlarms, + fetchUmeNe, + fetchUmeSyncStatus, + fetchUmeTokenStatus, + refreshUmeToken, +} from "../services/api"; +import { formatSystemTime } from "../utils/time"; + +function toIsoNoMs(d: Date): string { + return d.toISOString().replace(/\.\d{3}Z$/, "Z"); +} + +export function UmePage() { + const queryClient = useQueryClient(); + const [syncPage, setSyncPage] = useState(1); + const [syncPageSize, setSyncPageSize] = useState(20); + const [neKeyword, setNeKeyword] = useState(""); + const [nePage, setNePage] = useState(1); + const [nePageSize, setNePageSize] = useState(50); + + const [curSeverity, setCurSeverity] = useState(""); + const [curCleared, setCurCleared] = useState(""); + const [curNeId, setCurNeId] = useState(""); + const [curKeyword, setCurKeyword] = useState(""); + const [curPage, setCurPage] = useState(1); + const [curPageSize, setCurPageSize] = useState(50); + + const [hisSeverity, setHisSeverity] = useState(""); + const [hisNeId, setHisNeId] = useState(""); + const [hisKeyword, setHisKeyword] = useState(""); + const [hisTimeFrom, setHisTimeFrom] = useState(""); + const [hisTimeTo, setHisTimeTo] = useState(""); + const [hisPage, setHisPage] = useState(1); + const [hisPageSize, setHisPageSize] = useState(50); + + const applyQuickRange = (hours: number) => { + const now = new Date(); + const from = new Date(now.getTime() - hours * 3600 * 1000); + setHisTimeFrom(toIsoNoMs(from)); + setHisTimeTo(toIsoNoMs(now)); + setHisPage(1); + }; + + const syncMutation = useMutation({ + mutationFn: async (domains: string[]) => apiPost<{ ok: boolean; jobs: unknown[] }>("/v1/ume/sync", { domains }), + onSuccess: async () => { + await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] }); + await queryClient.invalidateQueries({ queryKey: ["umeNE"] }); + await queryClient.invalidateQueries({ queryKey: ["umeCurrentAlarms"] }); + await queryClient.invalidateQueries({ queryKey: ["umeHistoryAlarms"] }); + }, + }); + + const syncStatusQuery = useQuery({ + queryKey: ["umeSyncStatus", syncPage, syncPageSize], + queryFn: () => fetchUmeSyncStatus({ page: syncPage, pageSize: syncPageSize }), + staleTime: 5000, + }); + const tokenStatusQuery = useQuery({ + queryKey: ["umeTokenStatus"], + queryFn: fetchUmeTokenStatus, + staleTime: 3000, + refetchInterval: 5000, + }); + const tokenRefreshMutation = useMutation({ + mutationFn: refreshUmeToken, + onSuccess: async () => { + await queryClient.invalidateQueries({ queryKey: ["umeTokenStatus"] }); + }, + }); + const tokenDisconnectMutation = useMutation({ + mutationFn: disconnectUmeToken, + onSuccess: async () => { + await queryClient.invalidateQueries({ queryKey: ["umeTokenStatus"] }); + }, + }); + + const tokenExpiresIn = Number(tokenStatusQuery.data?.expires_in_s || 0); + const tokenLevel = !tokenStatusQuery.data?.has_token + ? "down" + : tokenExpiresIn < 15 + ? "down" + : tokenExpiresIn < 60 + ? "unknown" + : "up"; + const neQuery = useQuery({ + queryKey: ["umeNE", neKeyword, nePage, nePageSize], + queryFn: () => fetchUmeNe({ keyword: neKeyword, page: nePage, pageSize: nePageSize }), + staleTime: 5000, + }); + const currentQuery = useQuery({ + queryKey: ["umeCurrentAlarms", curSeverity, curCleared, curNeId, curKeyword, curPage, curPageSize], + queryFn: () => + fetchUmeCurrentAlarms({ + severity: curSeverity, + isCleared: curCleared, + neId: curNeId, + keyword: curKeyword, + page: curPage, + pageSize: curPageSize, + }), + staleTime: 5000, + }); + const historyQuery = useQuery({ + queryKey: ["umeHistoryAlarms", hisSeverity, hisNeId, hisKeyword, hisTimeFrom, hisTimeTo, hisPage, hisPageSize], + queryFn: () => + fetchUmeHistoryAlarms({ + severity: hisSeverity, + neId: hisNeId, + keyword: hisKeyword, + timeFrom: hisTimeFrom, + timeTo: hisTimeTo, + page: hisPage, + pageSize: hisPageSize, + }), + staleTime: 5000, + }); + + return ( + <> +
+
+

UME Token 状态

+
+ + token: {tokenStatusQuery.data?.has_token ? "connected" : "disconnected"} + + + expires_in: {typeof tokenStatusQuery.data?.expires_in_s === "number" ? `${tokenStatusQuery.data.expires_in_s}s` : "-"} + + {tokenStatusQuery.data?.token_preview ? preview: {tokenStatusQuery.data.token_preview} : null} +
+
+ + + +
+ {(tokenRefreshMutation.error || tokenDisconnectMutation.error) && ( +
+ 操作失败: {String(tokenRefreshMutation.error || tokenDisconnectMutation.error)} +
+ )} +
+
+

UME 同步

+
+ + + + +
+ {syncMutation.error &&
同步失败: {String(syncMutation.error)}
} +
+
+ +
+

同步状态

+
+ +
+ + + + + + + + + + + + + + + + {(syncStatusQuery.data?.items || []).map((x) => ( + + + + + + + + + + + + ))} + {!syncStatusQuery.isLoading && (syncStatusQuery.data?.items || []).length === 0 && ( + + + + )} + +
IDdomainstatuspulledinsertedupdatedstarted_atended_aterror
{x.id}{x.domain}{x.status}{x.pulled_count}{x.inserted_count}{x.updated_count}{formatSystemTime(x.started_at)}{x.ended_at ? formatSystemTime(x.ended_at) : "-"}{x.error_message || "-"}
暂无同步记录
+
+
+ 共 {syncStatusQuery.data?.total || 0} 条 · 第 {syncPage}/ + {Math.max(1, Math.ceil(Math.max(0, Number(syncStatusQuery.data?.total || 0)) / Math.max(1, syncPageSize)))} 页 +
+
+ + + +
+
+
+ +
+

网元清单

+
+ setNeKeyword(e.target.value)} /> + +
+ + + + + + + + + + + + + {(neQuery.data?.items || []).map((x) => ( + + + + + + + + + ))} + +
ne_iduser_labelne_nameiptypelast_seen
{x.ne_id}{x.user_label}{x.ne_name}{x.ip_address}{x.ne_type}{x.last_seen_at ? formatSystemTime(x.last_seen_at) : "-"}
+
+
+ 共 {neQuery.data?.total || 0} 条 · 第 {nePage}/ + {Math.max(1, Math.ceil(Math.max(0, Number(neQuery.data?.total || 0)) / Math.max(1, nePageSize)))} 页 +
+
+ + + +
+
+
+ +
+

当前告警

+
+ setCurKeyword(e.target.value)} /> + setCurNeId(e.target.value)} /> + + +
+ + + + + + + + + + + + {(currentQuery.data?.items || []).map((x) => ( + + + + + + + + ))} + +
time_createdseverityuser_labelobject_namecause
{x.time_created}{x.perceived_severity} + + {x.object_name}{x.native_probable_cause}
+
+
共 {currentQuery.data?.total || 0} 条 · 第 {curPage} 页
+
+ + + +
+
+
+ +
+

历史告警

+
+ setHisKeyword(e.target.value)} /> + setHisNeId(e.target.value)} /> + setHisTimeFrom(e.target.value)} /> + setHisTimeTo(e.target.value)} /> + +
+
+ + + + +
+ + + + + + + + + + + + {(historyQuery.data?.items || []).map((x) => ( + + + + + + + + ))} + +
time_createdseverityuser_labelobject_namecause
{x.time_created}{x.perceived_severity} + + {x.object_name}{x.native_probable_cause}
+
+
共 {historyQuery.data?.total || 0} 条 · 第 {hisPage} 页
+
+ + + +
+
+
+ + ); +} diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 03036b4..d63dcaa 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -1,4 +1,15 @@ -import type { AiAnalyzeHistoryResponse, AlarmQueryResponse, Batch, Diagnostics, ImportHistoryItem, IntegrationStatus } from "../types"; +import type { + AiAnalyzeHistoryResponse, + AlarmQueryResponse, + Batch, + Diagnostics, + ImportHistoryItem, + IntegrationStatus, + UmeAlarmItem, + UmeNeItem, + UmeSyncStatusResponse, + UmeTokenStatus, +} from "../types"; export const apiGet = async (path: string): Promise => { const res = await fetch(path, { headers: { accept: "application/json" } }); @@ -89,3 +100,61 @@ export const fetchAlarms = (params: { return apiGet(`/v1/alarms?${p.toString()}`); }; +export const fetchUmeSyncStatus = (params: { page: number; pageSize: number }) => { + const p = new URLSearchParams(); + p.set("page", String(Math.max(1, Number(params.page || 1)))); + p.set("page_size", String(Math.max(1, Math.min(200, Number(params.pageSize || 20))))); + return apiGet(`/v1/ume/sync/status?${p.toString()}`); +}; +export const fetchUmeTokenStatus = () => apiGet("/v1/ume/token/status"); +export const refreshUmeToken = () => apiPost("/v1/ume/token/refresh", {}); +export const disconnectUmeToken = () => apiPost("/v1/ume/token/disconnect", {}); + +export const fetchUmeNe = (params: { keyword: string; page: number; pageSize: number }) => { + const p = new URLSearchParams(); + if (params.keyword) p.set("keyword", params.keyword); + p.set("page", String(Math.max(1, Number(params.page || 1)))); + p.set("page_size", String(Math.max(1, Math.min(500, Number(params.pageSize || 50))))); + return apiGet<{ total: number; page: number; page_size: number; items: UmeNeItem[] }>(`/v1/ume/inventory/ne?${p.toString()}`); +}; + +export const fetchUmeCurrentAlarms = (params: { + severity: string; + isCleared: string; + neId: string; + keyword: string; + page: number; + pageSize: number; +}) => { + const p = new URLSearchParams(); + if (params.severity) p.set("severity", params.severity); + if (params.isCleared) p.set("is_cleared", params.isCleared); + if (params.neId) p.set("ne_id", params.neId); + if (params.keyword) p.set("keyword", params.keyword); + p.set("page", String(Math.max(1, Number(params.page || 1)))); + p.set("page_size", String(Math.max(1, Math.min(500, Number(params.pageSize || 50))))); + return apiGet<{ total: number; page: number; page_size: number; items: UmeAlarmItem[] }>(`/v1/ume/alarms?${p.toString()}`); +}; + +export const fetchUmeHistoryAlarms = (params: { + severity: string; + neId: string; + keyword: string; + timeFrom: string; + timeTo: string; + page: number; + pageSize: number; +}) => { + const p = new URLSearchParams(); + if (params.severity) p.set("severity", params.severity); + if (params.neId) p.set("ne_id", params.neId); + if (params.keyword) p.set("keyword", params.keyword); + if (params.timeFrom) p.set("time_from", params.timeFrom); + if (params.timeTo) p.set("time_to", params.timeTo); + p.set("page", String(Math.max(1, Number(params.page || 1)))); + p.set("page_size", String(Math.max(1, Math.min(500, Number(params.pageSize || 50))))); + return apiGet<{ total: number; page: number; page_size: number; items: UmeAlarmItem[] }>( + `/v1/ume/alarms/history?${p.toString()}`, + ); +}; + diff --git a/web/src/types.ts b/web/src/types.ts index d26230e..edf6c53 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -73,3 +73,58 @@ export type AiAnalyzeHistoryResponse = { page_size: number; items: AiAnalyzeHistoryItem[]; }; + +export type UmeSyncJobItem = { + id: number; + domain: string; + status: string; + trigger_mode: string; + pulled_count: number; + inserted_count: number; + updated_count: number; + error_message?: string; + started_at: string; + ended_at?: string | null; +}; + +export type UmeSyncStatusResponse = { + total?: number; + page?: number; + page_size?: number; + items: UmeSyncJobItem[]; + latest_by_domain?: Record; +}; + +export type UmeNeItem = { + ne_id: string; + ne_name: string; + user_label: string; + ip_address: string; + ne_type: string; + last_seen_at?: string; +}; + +export type UmeAlarmItem = { + alarm_key: string; + ne_id: string; + ne_name: string; + user_label: string; + object_name: string; + event_type: string; + native_probable_cause: string; + perceived_severity: string; + is_cleared: string; + time_created: string; + last_seen_at?: string; +}; + +export type UmeTokenStatus = { + ok: boolean; + has_token: boolean; + expires_in_s: number; + expires_at_epoch_s: number; + auth_header: string; + token_preview?: string; + error_kind?: string; + error?: string; +};