diff --git a/netx_api/config.py b/netx_api/config.py index ca2c12c..58da2fb 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -42,12 +42,13 @@ class Settings(BaseSettings): ume_keepalive_renew_before_s: int = 900 ume_sync_alarms_current_enabled: bool = True ume_sync_alarms_current_interval_s: int = 300 + ume_sync_inventory_auto_enabled: bool = True + ume_sync_inventory_every_hours: int = 48 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 diff --git a/netx_api/main.py b/netx_api/main.py index 546d370..5082995 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -70,7 +70,10 @@ _SQL_FORBIDDEN_RE = re.compile( _UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = { "token_keepalive": {"task": "token_keepalive", "status": "init", "last_run_at": None, "last_error": ""}, "alarms_current_auto_sync": {"task": "alarms_current_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, + "inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""}, } +_UME_RUNTIME_PAUSED: dict[str, bool] = {} +UME_KNOWN_RUNTIME_TASKS: tuple[str, ...] = tuple(_UME_RUNTIME_TASKS.keys()) _UME_RUNTIME_LOCK = threading.Lock() @@ -84,15 +87,38 @@ def _set_runtime_task(task: str, *, status: str, last_run_at: datetime | None = _UME_RUNTIME_TASKS[task] = item +def _runtime_is_paused(task: str) -> bool: + with _UME_RUNTIME_LOCK: + return bool(_UME_RUNTIME_PAUSED.get(str(task or "").strip())) + + +def _runtime_pause_task(task: str) -> None: + tid = str(task or "").strip() + with _UME_RUNTIME_LOCK: + if tid not in _UME_RUNTIME_TASKS: + raise KeyError(tid) + _UME_RUNTIME_PAUSED[tid] = True + + +def _runtime_resume_task(task: str) -> None: + tid = str(task or "").strip() + with _UME_RUNTIME_LOCK: + _UME_RUNTIME_PAUSED[tid] = False + + def _list_runtime_tasks() -> list[dict[str, Any]]: with _UME_RUNTIME_LOCK: out: list[dict[str, Any]] = [] for v in _UME_RUNTIME_TASKS.values(): + task_id = str(v.get("task") or "") + paused = bool(_UME_RUNTIME_PAUSED.get(task_id)) + eff_status = "paused" if paused else str(v.get("status") or "unknown") ts = _ensure_utc(v.get("last_run_at")) if isinstance(v.get("last_run_at"), datetime) else None out.append( { - "task": str(v.get("task") or ""), - "status": str(v.get("status") or "unknown"), + "task": task_id, + "status": eff_status, + "paused": paused, "last_run_at": ts.isoformat() if ts else None, "last_error": str(v.get("last_error") or ""), } @@ -353,6 +379,9 @@ def on_startup() -> None: # Best-effort keepalive: if token exists, periodically handshake to extend TTL. while True: try: + if _runtime_is_paused("token_keepalive"): + time.sleep(1) + continue client = _ume_client() st = client.token_status() expires_in = int(st.get("expires_in_s") or 0) @@ -375,6 +404,9 @@ def on_startup() -> None: def _alarms_current_sync_loop() -> None: while True: try: + if _runtime_is_paused("alarms_current_auto_sync"): + time.sleep(1) + continue db = SessionLocal() try: client = _ume_client() @@ -400,6 +432,43 @@ def on_startup() -> None: t2.start() except Exception: pass + try: + if bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)): + hours = int(getattr(settings, "ume_sync_inventory_every_hours", 48) or 48) + hours = max(1, min(hours, 168)) + interval_s = int(hours * 3600) + + def _inventory_auto_sync_loop() -> None: + while True: + try: + if _runtime_is_paused("inventory_auto_sync"): + time.sleep(1) + continue + db = SessionLocal() + try: + client = _ume_client() + sync_inventory_full(db, client, trigger_mode="schedule") + _set_runtime_task( + "inventory_auto_sync", + status="running", + last_run_at=datetime.now(timezone.utc), + last_error="", + ) + finally: + db.close() + except Exception as exc: + _set_runtime_task( + "inventory_auto_sync", + status="error", + last_run_at=datetime.now(timezone.utc), + last_error=str(exc)[:240], + ) + time.sleep(interval_s) + + t3 = threading.Thread(target=_inventory_auto_sync_loop, name="ume-inventory-auto-sync", daemon=True) + t3.start() + except Exception: + pass @app.get("/health") @@ -564,6 +633,26 @@ def ume_sync_status( } +@app.post("/v1/ume/runtime/tasks/{task}/pause") +def ume_runtime_task_pause(task: str) -> dict[str, Any]: + tid = str(task or "").strip() + if tid not in UME_KNOWN_RUNTIME_TASKS: + raise HTTPException(status_code=404, detail="unknown_runtime_task") + _runtime_pause_task(tid) + _set_runtime_task(tid, status="paused", last_error="") + return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()} + + +@app.post("/v1/ume/runtime/tasks/{task}/resume") +def ume_runtime_task_resume(task: str) -> dict[str, Any]: + tid = str(task or "").strip() + if tid not in UME_KNOWN_RUNTIME_TASKS: + raise HTTPException(status_code=404, detail="unknown_runtime_task") + _runtime_resume_task(tid) + _set_runtime_task(tid, status="running", last_error="") + return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()} + + @app.get("/v1/ume/inventory/ne") def ume_list_ne( keyword: str | None = Query(default=None), diff --git a/web/src/pages/UmePage.tsx b/web/src/pages/UmePage.tsx index 89328bc..6ed10dd 100644 --- a/web/src/pages/UmePage.tsx +++ b/web/src/pages/UmePage.tsx @@ -107,6 +107,17 @@ export function UmePage() { const runningTasks = (syncStatusQuery.data?.items || []).filter((x) => String(x.status || "").toLowerCase() === "running"); const runtimeTasks = syncStatusQuery.data?.runtime_tasks || []; + const runtimeTaskMutation = useMutation({ + mutationFn: async (vars: { task: string; action: "pause" | "resume" }) => + apiPost<{ ok: boolean }>( + `/v1/ume/runtime/tasks/${encodeURIComponent(vars.task)}/${vars.action}`, + {}, + ), + onSuccess: async () => { + await queryClient.invalidateQueries({ queryKey: ["umeSyncStatus"] }); + }, + }); + return ( <>
@@ -205,6 +216,7 @@ export function UmePage() { status last_run_at last_error + 操作 @@ -214,11 +226,32 @@ export function UmePage() { {x.status} {x.last_run_at ? formatSystemTime(x.last_run_at) : "-"} {x.last_error || "-"} + + {Boolean(x.paused) ? ( + + ) : ( + + )} + ))} {!syncStatusQuery.isLoading && runtimeTasks.length === 0 && ( - 暂无后台任务状态 + 暂无后台任务状态 )} @@ -462,7 +495,7 @@ export function UmePage() { {(currentQuery.data?.items || []).map((x) => ( - {x.time_created} + {formatSystemTime(x.time_created)} {x.perceived_severity} {x.ne_id} diff --git a/web/src/types.ts b/web/src/types.ts index 6660c2f..936b926 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -98,6 +98,7 @@ export type UmeSyncStatusResponse = { runtime_tasks?: Array<{ task: string; status: string; + paused?: boolean; last_run_at?: string | null; last_error?: string; }>;