UME: system time for alarm time_created, 48h inventory auto-sync, runtime pause/resume

- formatSystemTime for current alarm time_created\n- Background thread inventory_auto_sync (NETX_UME_SYNC_INVENTORY_AUTO_ENABLED, NETX_UME_SYNC_INVENTORY_EVERY_HOURS default 48)\n- Pause/resume for token_keepalive, alarms_current_auto_sync, inventory_auto_sync via POST /v1/ume/runtime/tasks/{task}/pause|resume\n- UmePage后台任务列操作 暂停/开始

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-05-11 19:18:28 +08:00
parent edad7b24fa
commit 1741965de6
4 changed files with 129 additions and 5 deletions

View file

@ -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

View file

@ -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),

View file

@ -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 (
<>
<section className="cards">
@ -205,6 +216,7 @@ export function UmePage() {
<th>status</th>
<th>last_run_at</th>
<th>last_error</th>
<th>操作</th>
</tr>
</thead>
<tbody>
@ -214,11 +226,32 @@ export function UmePage() {
<td>{x.status}</td>
<td>{x.last_run_at ? formatSystemTime(x.last_run_at) : "-"}</td>
<td>{x.last_error || "-"}</td>
<td>
{Boolean(x.paused) ? (
<button
type="button"
className="link-btn"
disabled={runtimeTaskMutation.isPending}
onClick={() => runtimeTaskMutation.mutate({ task: x.task, action: "resume" })}
>
开始
</button>
) : (
<button
type="button"
className="link-btn"
disabled={runtimeTaskMutation.isPending}
onClick={() => runtimeTaskMutation.mutate({ task: x.task, action: "pause" })}
>
暂停
</button>
)}
</td>
</tr>
))}
{!syncStatusQuery.isLoading && runtimeTasks.length === 0 && (
<tr>
<td colSpan={4}>暂无后台任务状态</td>
<td colSpan={5}>暂无后台任务状态</td>
</tr>
)}
</tbody>
@ -462,7 +495,7 @@ export function UmePage() {
<tbody>
{(currentQuery.data?.items || []).map((x) => (
<tr key={x.alarm_key}>
<td>{x.time_created}</td>
<td>{formatSystemTime(x.time_created)}</td>
<td>{x.perceived_severity}</td>
<td>{x.ne_id}</td>
<td>

View file

@ -98,6 +98,7 @@ export type UmeSyncStatusResponse = {
runtime_tasks?: Array<{
task: string;
status: string;
paused?: boolean;
last_run_at?: string | null;
last_error?: string;
}>;