mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 04:20:45 +08:00
feat(UME): 增加RESTCONF对接与可视化运维页
- 支持 token 获取/续约/断开、DB共享缓存与跨进程单飞锁、keepalive 保活 - 新增 inventory/current/history 告警同步入库与查询接口,前端增加 UME 对接页面与分页 - 修复 stop_netx.ps1 在部分 PowerShell 版本下 Stop-Job -Force 报错 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
ab1b76dd47
commit
7f2e4ec393
15 changed files with 2577 additions and 15 deletions
|
|
@ -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()
|
||||
|
|
|
|||
390
netx_api/main.py
390
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.
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
458
netx_api/ume_client.py
Normal file
458
netx_api/ume_client.py
Normal file
|
|
@ -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
|
||||
299
netx_api/ume_sync_service.py
Normal file
299
netx_api/ume_sync_service.py
Normal file
|
|
@ -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)
|
||||
143
netx_api/ume_token_store.py
Normal file
143
netx_api/ume_token_store.py
Normal file
|
|
@ -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()
|
||||
Loading…
Add table
Add a link
Reference in a new issue