mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-09 04:40:45 +08:00
831 lines
33 KiB
Python
831 lines
33 KiB
Python
"""netx ops helpers: runtime context inject + optional legacy inline HTTP tools.
|
||
|
||
Default: tools are exposed via stdio MCP (``pip install -e packages/netx-mcp`` then ``python -m netx_mcp``).
|
||
Set ``OCLAW_NETX_BUILTIN_TOOLS=1`` to re-register inline ``netx_*`` expert tools.
|
||
|
||
Configure via environment:
|
||
- ``OCLAW_NETX_BASE_URL`` (default ``http://127.0.0.1:8890``) — runtime anchor HTTP probe
|
||
- ``OCLAW_NETX_API_TOKEN`` (optional) → sent as ``Authorization: Bearer …`` if set.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import contextvars
|
||
import os
|
||
import threading
|
||
import time
|
||
from typing import Any
|
||
|
||
import httpx
|
||
|
||
from runtime.tools.base import ToolSpec
|
||
|
||
# Set by ToolExecutor for netx_* tools so responses match session language.
|
||
NETX_TOOL_LANG: contextvars.ContextVar[str] = contextvars.ContextVar("netx_tool_lang", default="zh")
|
||
|
||
_PROTOCOL_KEY_ZH_TO_EN: dict[str, str] = {
|
||
"其他": "Other",
|
||
"时钟": "Clock",
|
||
"OTN/光": "OTN/Optical",
|
||
"电源": "Power",
|
||
}
|
||
|
||
|
||
_UME_RAW_GROUP_FIELDS = [
|
||
"alarm_alarm_key",
|
||
"alarm_host_name",
|
||
"alarm_ne_id",
|
||
"alarm_object_name",
|
||
"alarm_event_type",
|
||
"alarm_native_probable_cause",
|
||
"alarm_perceived_severity",
|
||
"alarm_is_cleared",
|
||
"alarm_time_created",
|
||
"alarm_root_cause_alarm_indication",
|
||
"ne_ne_id",
|
||
"ne_ne_name",
|
||
"ne_user_label",
|
||
"ne_ip_address",
|
||
"ne_ipv6_address",
|
||
"ne_ne_type",
|
||
"ne_device_level",
|
||
"ne_host_name",
|
||
"ne_location",
|
||
"ne_hardware_version",
|
||
"ne_loopback",
|
||
"ne_consistent_state",
|
||
"ne_interface_version",
|
||
"ne_mac",
|
||
"ne_admin_status",
|
||
"ne_address_type",
|
||
"ne_connection_status",
|
||
"ne_maintain_status",
|
||
"ne_net_mask",
|
||
"ne_create_time",
|
||
"ne_creator",
|
||
"ne_vendor",
|
||
"ne_source_type",
|
||
"ne_exists",
|
||
]
|
||
|
||
|
||
def _netx_base_url() -> str:
|
||
return (os.getenv("OCLAW_NETX_BASE_URL") or "http://127.0.0.1:8890").strip().rstrip("/")
|
||
|
||
|
||
def _netx_headers() -> dict[str, str]:
|
||
h = {"accept": "application/json"}
|
||
tok = (os.getenv("OCLAW_NETX_API_TOKEN") or "").strip()
|
||
if tok:
|
||
h["authorization"] = f"Bearer {tok}"
|
||
return h
|
||
|
||
|
||
def _netx_lang_query_params() -> dict[str, str]:
|
||
lang = str(NETX_TOOL_LANG.get() or "zh").strip().lower()
|
||
if lang.startswith("en"):
|
||
return {"lang": "en"}
|
||
return {}
|
||
|
||
|
||
def _localize_netx_payload(data: dict[str, Any], *, lang: str) -> dict[str, Any]:
|
||
"""Map legacy Chinese protocol bucket labels to English for en sessions."""
|
||
if not str(lang or "").strip().lower().startswith("en"):
|
||
return data
|
||
proto = data.get("protocol_summary")
|
||
if isinstance(proto, list):
|
||
for row in proto:
|
||
if isinstance(row, dict):
|
||
k = str(row.get("key") or "")
|
||
if k in _PROTOCOL_KEY_ZH_TO_EN:
|
||
row["key"] = _PROTOCOL_KEY_ZH_TO_EN[k]
|
||
return data
|
||
|
||
|
||
def _http_post_json(path: str, body: dict[str, Any], *, timeout: float = 180.0) -> dict[str, Any]:
|
||
base = _netx_base_url()
|
||
url = f"{base}{path}"
|
||
try:
|
||
with httpx.Client(timeout=timeout, trust_env=False) as client:
|
||
resp = client.post(url, json=body, headers=_netx_headers())
|
||
text = resp.text
|
||
if not resp.is_success:
|
||
return {"ok": False, "error": f"netx_http_{resp.status_code}", "detail": text[:800]}
|
||
data = resp.json() if text else {}
|
||
if isinstance(data, dict):
|
||
data = _localize_netx_payload(data, lang=str(NETX_TOOL_LANG.get() or "zh"))
|
||
return {"ok": True, "data": data if isinstance(data, dict) else {"raw": data}}
|
||
except Exception as exc:
|
||
return {"ok": False, "error": "netx_request_failed", "detail": str(exc)[:800]}
|
||
|
||
|
||
def _http_json(method: str, path: str, *, params: dict[str, Any] | None = None) -> dict[str, Any]:
|
||
base = _netx_base_url()
|
||
url = f"{base}{path}"
|
||
merged: dict[str, Any] = dict(_netx_lang_query_params())
|
||
if params:
|
||
merged.update(params)
|
||
try:
|
||
# Do not inherit system proxy settings for local netx calls.
|
||
with httpx.Client(timeout=45.0, trust_env=False) as client:
|
||
resp = client.request(method, url, params=merged or None, headers=_netx_headers())
|
||
text = resp.text
|
||
if not resp.is_success:
|
||
return {"ok": False, "error": f"netx_http_{resp.status_code}", "detail": text[:800]}
|
||
data = resp.json() if text else {}
|
||
if isinstance(data, dict):
|
||
data = _localize_netx_payload(data, lang=str(NETX_TOOL_LANG.get() or "zh"))
|
||
return {"ok": True, "data": data if isinstance(data, dict) else {"raw": data}}
|
||
except Exception as exc:
|
||
return {"ok": False, "error": "netx_request_failed", "detail": str(exc)[:800]}
|
||
|
||
|
||
def _resolve_ume_anchor() -> dict[str, Any]:
|
||
"""Resolve current UME alarm anchor from netx sync status."""
|
||
r = _http_json("GET", "/v1/ume/sync/status", params={"page": 1, "page_size": 20})
|
||
if not r.get("ok"):
|
||
return {"ok": False, "error": "netx_ume_sync_status_failed", "detail": r.get("detail"), "upstream": r}
|
||
data = r.get("data") or {}
|
||
latest = data.get("latest_by_domain") if isinstance(data.get("latest_by_domain"), dict) else {}
|
||
cur = latest.get("alarms_current") if isinstance(latest.get("alarms_current"), dict) else {}
|
||
return {
|
||
"ok": True,
|
||
"anchor": {
|
||
"domain": "alarms_current",
|
||
"status": str(cur.get("status") or ""),
|
||
"trigger_mode": str(cur.get("trigger_mode") or ""),
|
||
"started_at": str(cur.get("started_at") or ""),
|
||
"ended_at": str(cur.get("ended_at") or ""),
|
||
"pulled_count": int(cur.get("pulled_count") or 0),
|
||
"inserted_count": int(cur.get("inserted_count") or 0),
|
||
"updated_count": int(cur.get("updated_count") or 0),
|
||
"error_message": str(cur.get("error_message") or ""),
|
||
},
|
||
}
|
||
|
||
|
||
_OPS_NETX_SYS_CTX_LOCK = threading.Lock()
|
||
# Lang code -> (monotonic_ts, formatted extension text); short TTL to avoid hammering netx each tool round.
|
||
_OPS_NETX_SYS_CTX_CACHE: dict[str, tuple[float, str]] = {}
|
||
_OPS_NETX_SYS_CTX_TTL_SEC = 5.0
|
||
|
||
|
||
def _format_ops_netx_system_extension(r: dict[str, Any], *, lang_en: bool) -> str:
|
||
if r.get("ok"):
|
||
row = r.get("anchor") if isinstance(r.get("anchor"), dict) else {}
|
||
status = str(row.get("status") or "")
|
||
mode = str(row.get("trigger_mode") or "")
|
||
started = str(row.get("started_at") or "")
|
||
ended = str(row.get("ended_at") or "")
|
||
pulled = int(row.get("pulled_count") or 0)
|
||
inserted = int(row.get("inserted_count") or 0)
|
||
updated = int(row.get("updated_count") or 0)
|
||
err = str(row.get("error_message") or "").strip()
|
||
lines_en = [
|
||
"[Netx UME current-alarms anchor]",
|
||
f"- status: {status}",
|
||
f"- trigger_mode: {mode}",
|
||
f"- started_at: {started}",
|
||
f"- ended_at: {ended}",
|
||
]
|
||
lines_zh = [
|
||
"[当前 netx UME告警锚点]",
|
||
f"- 状态: {status}",
|
||
f"- 触发方式: {mode}",
|
||
f"- 开始时间: {started}",
|
||
f"- 结束时间: {ended}",
|
||
]
|
||
(lines_en if lang_en else lines_zh).append(
|
||
f"- pulled/inserted/updated: {pulled}/{inserted}/{updated}"
|
||
if lang_en
|
||
else f"- 拉取/新增/更新: {pulled}/{inserted}/{updated}"
|
||
)
|
||
if err:
|
||
(lines_en if lang_en else lines_zh).append(
|
||
f"- last_error: {err[:200]}" if lang_en else f"- 最近错误: {err[:200]}"
|
||
)
|
||
tail_en = (
|
||
"- MCP tools (server_id=netx): mcp__netx__queryUmeAlarms, mcp__netx__aggregateUmeAlarms, "
|
||
"mcp__netx__runUmeDiagnostics, mcp__netx__queryUmeNeInventory, mcp__netx__getUmeNe, "
|
||
"mcp__netx__queryUmeAlarmsRaw, mcp__netx__aggregateUmeAlarmsRaw, mcp__netx__listUmeAlarmFields, "
|
||
"mcp__netx__sqlQueryUme, mcp__netx__listManagedNe, mcp__netx__getManagedNe, mcp__netx__execManagedNe\n"
|
||
"- note: this is only runtime anchor; use tools for alarm/ne evidence.\n"
|
||
"- English session: user-visible reply must contain NO Chinese/CJK; translate alarm text fields."
|
||
)
|
||
tail_zh = (
|
||
"- MCP 工具(server_id=netx):mcp__netx__queryUmeAlarms、mcp__netx__aggregateUmeAlarms、"
|
||
"mcp__netx__runUmeDiagnostics、mcp__netx__queryUmeNeInventory、mcp__netx__getUmeNe、"
|
||
"mcp__netx__queryUmeAlarmsRaw、mcp__netx__aggregateUmeAlarmsRaw、mcp__netx__listUmeAlarmFields、"
|
||
"mcp__netx__sqlQueryUme、mcp__netx__listManagedNe、mcp__netx__getManagedNe、mcp__netx__execManagedNe\n"
|
||
"- 说明: 此处仅为运行锚点;具体告警/网元信息必须以工具返回为准,勿臆测。"
|
||
)
|
||
return "\n".join(lines_en + [tail_en]) if lang_en else "\n".join(lines_zh + [tail_zh])
|
||
err = str(r.get("error") or "")
|
||
detail = str(r.get("detail") or "")[:240]
|
||
if lang_en:
|
||
return (
|
||
"[Netx UME current-alarms anchor]\n"
|
||
f"- error: {err}\n"
|
||
f"- detail: {detail}\n"
|
||
"- fix: check NETX_API_URL (or OCLAW_NETX_BASE_URL) and that netx API is reachable; "
|
||
"ensure MCP server_id=netx is bound and synced."
|
||
)
|
||
return (
|
||
"[当前 netx UME告警锚点]\n"
|
||
f"- 错误: {err}\n"
|
||
f"- 详情: {detail}\n"
|
||
"- 处理: 检查 NETX_API_URL(或 OCLAW_NETX_BASE_URL)与 netx 服务是否可达;"
|
||
"确认 Admin 已安装并绑定 MCP server_id=netx。"
|
||
)
|
||
|
||
|
||
def ops_netx_system_context_extension(*, lang: str = "zh") -> str:
|
||
"""Append to ops specialist system prompt: UME sync anchor (direct_loop injection).
|
||
|
||
Cached briefly to reduce duplicate HTTP calls across tool rounds.
|
||
"""
|
||
if str(os.getenv("OCLAW_OPS_NETX_CONTEXT_INJECT") or "1").strip().lower() in {"0", "false", "no", "off"}:
|
||
return ""
|
||
lang_en = str(lang or "").strip().lower().startswith("en")
|
||
lk = "en" if lang_en else "zh"
|
||
now = time.monotonic()
|
||
with _OPS_NETX_SYS_CTX_LOCK:
|
||
hit = _OPS_NETX_SYS_CTX_CACHE.get(lk)
|
||
if hit and (now - hit[0]) < _OPS_NETX_SYS_CTX_TTL_SEC:
|
||
return hit[1]
|
||
r = _resolve_ume_anchor()
|
||
text = _format_ops_netx_system_extension(r, lang_en=lang_en)
|
||
store_ts = time.monotonic()
|
||
with _OPS_NETX_SYS_CTX_LOCK:
|
||
_OPS_NETX_SYS_CTX_CACHE[lk] = (store_ts, text)
|
||
return text
|
||
|
||
|
||
def netx_query_ume_alarms_tool() -> ToolSpec:
|
||
"""Paginated UME current alarms from netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
# Guardrail: avoid accidental full scans by endless paging.
|
||
page = max(1, int(args.get("page") or 1))
|
||
if page > 2:
|
||
page = 2
|
||
page_size = min(500, max(1, int(args.get("page_size") or 50)))
|
||
params: dict[str, Any] = {"page": page, "page_size": page_size}
|
||
if str(args.get("severity") or "").strip():
|
||
params["severity"] = str(args.get("severity")).strip()
|
||
ne_name = str(args.get("ne_name") or "").strip()
|
||
keyword = str(args.get("keyword") or "").strip()
|
||
if keyword:
|
||
params["keyword"] = keyword
|
||
elif ne_name:
|
||
params["keyword"] = ne_name
|
||
if str(args.get("ne_id") or "").strip():
|
||
params["ne_id"] = str(args.get("ne_id")).strip()
|
||
return _http_json("GET", "/v1/ume/alarms", params=params)
|
||
|
||
return ToolSpec(
|
||
name="netx_query_ume_alarms",
|
||
description=(
|
||
"读取 netx UME 当前告警明细(实时表);每条含 host_name(网元主展示键,同步时已写入告警表)。"
|
||
"支持 severity/ne_id/keyword 与分页。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"severity": {"type": "string"},
|
||
"ne_id": {"type": "string"},
|
||
"ne_name": {"type": "string", "description": "兼容参数,会映射到 keyword"},
|
||
"keyword": {"type": "string", "description": "按网元名/标签/IP/对象名等关键字检索"},
|
||
"page": {"type": "integer", "minimum": 1, "default": 1},
|
||
"page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50},
|
||
},
|
||
"required": [],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "alarms", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_aggregate_ume_alarms_tool() -> ToolSpec:
|
||
"""Aggregate UME current alarms from netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
_ = args
|
||
return _http_json("GET", "/v1/ume/alarms/aggregate", params=None)
|
||
|
||
return ToolSpec(
|
||
name="netx_aggregate_ume_alarms",
|
||
description="读取 netx UME 当前告警聚合(by_severity/by_ne)。",
|
||
parameters={"type": "object", "properties": {}, "required": [], "additionalProperties": False},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "alarms", "aggregate", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_run_ume_diagnostics_tool() -> ToolSpec:
|
||
"""Diagnostics summary for UME current alarms."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
_ = args
|
||
return _http_json("GET", "/v1/ume/diagnostics", params=None)
|
||
|
||
return ToolSpec(
|
||
name="netx_run_ume_diagnostics",
|
||
description="读取 netx UME 告警诊断摘要(级别分布、Top 告警码、Top 网元、协议归类)。",
|
||
parameters={"type": "object", "properties": {}, "required": [], "additionalProperties": False},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "diagnostics", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_query_ume_ne_inventory_tool() -> ToolSpec:
|
||
"""Paged UME NE inventory synced in netx (PostgreSQL-backed)."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
page = max(1, int(args.get("page") or 1))
|
||
page_size = min(500, max(1, int(args.get("page_size") or 50)))
|
||
params: dict[str, Any] = {"page": page, "page_size": page_size}
|
||
if str(args.get("keyword") or "").strip():
|
||
params["keyword"] = str(args.get("keyword")).strip()
|
||
return _http_json("GET", "/v1/ume/inventory/ne", params=params)
|
||
|
||
return ToolSpec(
|
||
name="netx_query_ume_ne_inventory",
|
||
description=(
|
||
"查询 netx 已同步的 UME 网元清单(读 /v1/ume/inventory/ne,与 netx Web「网元清单」同源)。"
|
||
"keyword 可选:匹配 ne_id / ne_name / user_label / ip_address / host_name(主机名)包含。"
|
||
"返回 total、page、page_size、items(含 host_name、在线状态、地址、类型等)。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"keyword": {"type": "string", "description": "关键字过滤(可选)"},
|
||
"page": {"type": "integer", "minimum": 1, "default": 1},
|
||
"page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50},
|
||
},
|
||
"required": [],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "inventory", "ne", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_get_ume_ne_tool() -> ToolSpec:
|
||
"""Single UME NE detail by ne_id from netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
from urllib.parse import quote
|
||
|
||
ne_id = str(args.get("ne_id") or "").strip()
|
||
if not ne_id:
|
||
return {"ok": False, "error": "ne_id_required", "error_code": "ne_id_required"}
|
||
safe = quote(ne_id, safe="")
|
||
return _http_json("GET", f"/v1/ume/inventory/ne/{safe}", params=None)
|
||
|
||
return ToolSpec(
|
||
name="netx_get_ume_ne",
|
||
description=(
|
||
"按网元 UUID(ne_id)读取 netx 中单条 UME 网元详情(GET /v1/ume/inventory/ne/{ne_id})。"
|
||
"含 vendor、source_type、raw_json 等;404 时上游返回 ume_ne_not_found。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"ne_id": {"type": "string", "description": "网元 UUID(与清单中 ne_id 一致)"},
|
||
},
|
||
"required": ["ne_id"],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "inventory", "ne", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_query_ume_alarms_raw_tool() -> ToolSpec:
|
||
"""Power query UME current alarms with full alarm+NE fields."""
|
||
|
||
presets: dict[str, list[str]] = {
|
||
"brief": [
|
||
"alarm_alarm_key",
|
||
"alarm_host_name",
|
||
"alarm_perceived_severity",
|
||
"alarm_event_type",
|
||
"alarm_last_seen_at",
|
||
"ne_host_name",
|
||
"ne_user_label",
|
||
"ne_ne_name",
|
||
"ne_ip_address",
|
||
"ne_exists",
|
||
],
|
||
"evidence": [
|
||
"alarm_alarm_key",
|
||
"alarm_host_name",
|
||
"alarm_object_name",
|
||
"alarm_event_type",
|
||
"alarm_native_probable_cause",
|
||
"alarm_perceived_severity",
|
||
"alarm_is_cleared",
|
||
"alarm_time_created",
|
||
"alarm_last_seen_at",
|
||
"ne_host_name",
|
||
"ne_user_label",
|
||
"ne_ne_name",
|
||
"ne_ip_address",
|
||
"ne_connection_status",
|
||
"ne_exists",
|
||
],
|
||
"ne_debug": [
|
||
"alarm_alarm_key",
|
||
"alarm_ne_id",
|
||
"alarm_perceived_severity",
|
||
"alarm_last_seen_at",
|
||
"ne_user_label",
|
||
"ne_ne_name",
|
||
"ne_ip_address",
|
||
"ne_ipv6_address",
|
||
"ne_device_level",
|
||
"ne_host_name",
|
||
"ne_connection_status",
|
||
"ne_admin_status",
|
||
"ne_address_type",
|
||
"ne_maintain_status",
|
||
"ne_exists",
|
||
],
|
||
}
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
page = max(1, int(args.get("page") or 1))
|
||
page_size = min(500, max(1, int(args.get("page_size") or 50)))
|
||
params: dict[str, Any] = {"page": page, "page_size": page_size}
|
||
for k in ("severity", "is_cleared", "ne_id", "event_type", "keyword", "time_from", "time_to", "order_by", "order"):
|
||
v = str(args.get(k) or "").strip()
|
||
if v:
|
||
params[k] = v
|
||
sf = args.get("select_fields")
|
||
fields: list[str] = []
|
||
if isinstance(sf, list):
|
||
fields = [str(x).strip() for x in sf if str(x).strip()]
|
||
if not fields:
|
||
preset = str(args.get("field_preset") or "").strip().lower()
|
||
fields = list(presets.get(preset) or [])
|
||
if fields:
|
||
params["select_fields"] = ",".join(fields)
|
||
return _http_json("GET", "/v1/ume/alarms/raw", params=params)
|
||
|
||
return ToolSpec(
|
||
name="netx_query_ume_alarms_raw",
|
||
description=(
|
||
"自由查询 netx UME 当前告警原始视图,返回 alarm_* + ne_* 全字段。"
|
||
"可按 severity/is_cleared/ne_id/event_type/keyword/time_from/time_to 过滤,支持排序分页。"
|
||
"select_fields 可按需指定返回字段,降低输出体积。"
|
||
"field_preset 可快速选用默认字段集(brief/evidence/ne_debug)。"
|
||
"建议先调用 netx_list_ume_alarm_fields 查看可用字段。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"severity": {"type": "string"},
|
||
"is_cleared": {"type": "string"},
|
||
"ne_id": {"type": "string"},
|
||
"event_type": {"type": "string"},
|
||
"keyword": {"type": "string"},
|
||
"time_from": {"type": "string", "description": "ISO8601 时间下界(按 last_seen_at)"},
|
||
"time_to": {"type": "string", "description": "ISO8601 时间上界(按 last_seen_at)"},
|
||
"order_by": {"type": "string", "enum": ["last_seen_at", "time_created", "perceived_severity", "event_type", "ne_id"]},
|
||
"order": {"type": "string", "enum": ["asc", "desc"]},
|
||
"select_fields": {
|
||
"type": "array",
|
||
"items": {"type": "string"},
|
||
"description": "可选返回字段,如 alarm_alarm_key/ne_user_label/ne_exists",
|
||
},
|
||
"field_preset": {
|
||
"type": "string",
|
||
"enum": ["brief", "evidence", "ne_debug"],
|
||
"description": "字段集预设;当未传 select_fields 时生效",
|
||
},
|
||
"page": {"type": "integer", "minimum": 1, "default": 1},
|
||
"page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50},
|
||
},
|
||
"required": [],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "alarms", "power_query", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_list_ume_alarm_fields_tool() -> ToolSpec:
|
||
"""List field names for UME raw alarm query."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
_ = args
|
||
return _http_json("GET", "/v1/ume/alarms/fields", params=None)
|
||
|
||
return ToolSpec(
|
||
name="netx_list_ume_alarm_fields",
|
||
description="列出 UME 当前告警 raw 查询可用字段(alarm_fields/ne_fields/order_by_allowed)。",
|
||
parameters={"type": "object", "properties": {}, "required": [], "additionalProperties": False},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "alarms", "schema", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_sql_query_ume_tool() -> ToolSpec:
|
||
"""Execute read-only SQL on UME tables in netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
sql = str(args.get("sql") or "").strip()
|
||
limit = max(1, min(2000, int(args.get("limit") or 200)))
|
||
statement_timeout_ms = max(0, min(30000, int(args.get("statement_timeout_ms") or 0)))
|
||
if not sql:
|
||
return {"ok": False, "error": "sql_required"}
|
||
base = _netx_base_url()
|
||
url = f"{base}/v1/sql/ume_query"
|
||
try:
|
||
# Do not inherit system proxy settings for local netx calls.
|
||
with httpx.Client(timeout=60.0, trust_env=False) as client:
|
||
resp = client.post(
|
||
url,
|
||
json={"sql": sql, "limit": limit, "statement_timeout_ms": statement_timeout_ms},
|
||
headers=_netx_headers(),
|
||
)
|
||
text = resp.text
|
||
if not resp.is_success:
|
||
return {"ok": False, "error": f"netx_http_{resp.status_code}", "detail": text[:800]}
|
||
data = resp.json() if text else {}
|
||
return {"ok": True, "data": data if isinstance(data, dict) else {"raw": data}}
|
||
except Exception as exc:
|
||
return {"ok": False, "error": "netx_request_failed", "detail": str(exc)[:800]}
|
||
|
||
return ToolSpec(
|
||
name="netx_sql_query_ume",
|
||
description=(
|
||
"在 netx 上执行 UME 只读 SQL(服务端强制 SELECT-only、单语句、限制表为 "
|
||
"ume_alarms_current/ume_inventory_ne,并强制 limit)。"
|
||
"推荐默认模板:设置 statement_timeout_ms=8000,且 SQL 带时间窗过滤(last_seen_at >= now() - interval '30 minutes')。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"sql": {"type": "string", "description": "只读 SELECT SQL;仅允许 UME 当前告警与网元表"},
|
||
"limit": {"type": "integer", "minimum": 1, "maximum": 2000, "default": 200},
|
||
"statement_timeout_ms": {
|
||
"type": "integer",
|
||
"minimum": 0,
|
||
"maximum": 30000,
|
||
"default": 0,
|
||
"description": "可选查询超时(ms);0 表示使用数据库默认超时",
|
||
},
|
||
},
|
||
"required": ["sql"],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "sql", "power_query", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_aggregate_ume_alarms_raw_tool() -> ToolSpec:
|
||
"""Dynamic aggregation on UME raw fields."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
params: dict[str, Any] = {}
|
||
for k in (
|
||
"group_by",
|
||
"group_by2",
|
||
"severity",
|
||
"is_cleared",
|
||
"ne_id",
|
||
"event_type",
|
||
"keyword",
|
||
"time_from",
|
||
"time_to",
|
||
"limit",
|
||
):
|
||
v = args.get(k)
|
||
if v is None:
|
||
continue
|
||
sv = str(v).strip()
|
||
if sv:
|
||
params[k] = sv
|
||
return _http_json("GET", "/v1/ume/alarms/aggregate/raw", params=params)
|
||
|
||
return ToolSpec(
|
||
name="netx_aggregate_ume_alarms_raw",
|
||
description=(
|
||
"按 UME raw 字段做动态聚合(group_by/group_by2),支持与 raw 同口径过滤条件。"
|
||
"group_by 需使用 alarm_*/ne_* 字段,建议先 netx_list_ume_alarm_fields。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"group_by": {
|
||
"type": "string",
|
||
"enum": _UME_RAW_GROUP_FIELDS,
|
||
"description": "主分组字段(网元主键优先 alarm_host_name 或 ne_host_name;勿用 alarm_ne_id/ne_ne_id)",
|
||
},
|
||
"group_by2": {"type": "string", "enum": _UME_RAW_GROUP_FIELDS, "description": "可选第二分组字段"},
|
||
"severity": {"type": "string"},
|
||
"is_cleared": {"type": "string"},
|
||
"ne_id": {"type": "string"},
|
||
"event_type": {"type": "string"},
|
||
"keyword": {"type": "string"},
|
||
"time_from": {"type": "string"},
|
||
"time_to": {"type": "string"},
|
||
"limit": {"type": "integer", "minimum": 1, "maximum": 2000, "default": 200},
|
||
},
|
||
"required": ["group_by"],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "ume", "alarms", "aggregate", "power_query", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_list_managed_ne_tool() -> ToolSpec:
|
||
"""List netx managed NEs (inventory for CLI login targets)."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
page = max(1, int(args.get("page") or 1))
|
||
page_size = min(500, max(1, int(args.get("page_size") or 50)))
|
||
params: dict[str, Any] = {"page": page, "page_size": page_size}
|
||
if str(args.get("keyword") or "").strip():
|
||
params["keyword"] = str(args.get("keyword")).strip()
|
||
if str(args.get("vendor") or "").strip():
|
||
params["vendor"] = str(args.get("vendor")).strip()
|
||
if str(args.get("connect_status") or "").strip():
|
||
params["connect_status"] = str(args.get("connect_status")).strip()
|
||
return _http_json("GET", "/v1/managed-ne", params=params)
|
||
|
||
return ToolSpec(
|
||
name="netx_list_managed_ne",
|
||
description=(
|
||
"列出 netx「网元管理」中已纳管的设备(GET /v1/managed-ne)。"
|
||
"返回 id、name、ip、vendor、device_type、connect_status 等(不含密码)。"
|
||
"keyword 可匹配名称/IP/用户名/标签;connect_status 可选 unknown/testing/pass/fail。"
|
||
"登录查配置前先用本工具定位 ne_id,再 netx_get_managed_ne / netx_exec_managed_ne。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"keyword": {"type": "string", "description": "名称/IP/用户名/标签包含(可选)"},
|
||
"vendor": {"type": "string", "description": "厂商过滤(可选)"},
|
||
"connect_status": {
|
||
"type": "string",
|
||
"enum": ["unknown", "testing", "pass", "fail"],
|
||
"description": "连通性状态过滤(可选)",
|
||
},
|
||
"page": {"type": "integer", "minimum": 1, "default": 1},
|
||
"page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50},
|
||
},
|
||
"required": [],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "managed_ne", "inventory", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_get_managed_ne_tool() -> ToolSpec:
|
||
"""Single managed NE metadata from netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
ne_id = str(args.get("ne_id") or "").strip()
|
||
if not ne_id:
|
||
return {"ok": False, "error": "ne_id_required", "error_code": "ne_id_required"}
|
||
return _http_json("GET", f"/v1/managed-ne/{ne_id}", params=None)
|
||
|
||
return ToolSpec(
|
||
name="netx_get_managed_ne",
|
||
description=(
|
||
"读取 netx 单条纳管网元详情(GET /v1/managed-ne/{ne_id})。"
|
||
"含 connect_status、connect_message、connect_detail(连通测试日志)、跳板 hop_* 配置摘要。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"ne_id": {"type": "string", "description": "网元 UUID(列表 items[].id)"},
|
||
},
|
||
"required": ["ne_id"],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "managed_ne", "read_only"}),
|
||
risk_level="low",
|
||
read_only=True,
|
||
)
|
||
|
||
|
||
def netx_exec_managed_ne_tool() -> ToolSpec:
|
||
"""Run read-only CLI on a managed NE via netx."""
|
||
|
||
def handler(args: dict[str, Any]) -> dict[str, Any]:
|
||
ne_id = str(args.get("ne_id") or "").strip()
|
||
if not ne_id:
|
||
return {"ok": False, "error": "ne_id_required", "error_code": "ne_id_required"}
|
||
raw_cmds = args.get("commands")
|
||
if not isinstance(raw_cmds, list) or not raw_cmds:
|
||
return {"ok": False, "error": "commands_required", "error_code": "commands_required"}
|
||
commands = [str(c).strip() for c in raw_cmds if str(c).strip()]
|
||
if not commands:
|
||
return {"ok": False, "error": "commands_required", "error_code": "commands_required"}
|
||
if len(commands) > 5:
|
||
return {"ok": False, "error": "too_many_commands", "error_code": "too_many_commands"}
|
||
body: dict[str, Any] = {"ne_id": ne_id, "commands": commands}
|
||
rts = args.get("read_timeout_sec")
|
||
if rts is not None:
|
||
body["read_timeout_sec"] = int(rts)
|
||
out = _http_post_json("/v1/managed-ne/exec", body, timeout=300.0)
|
||
if not out.get("ok"):
|
||
return out
|
||
data = out.get("data") or {}
|
||
if isinstance(data, dict) and data.get("ok") is False:
|
||
return {"ok": False, "data": data, "error": str(data.get("error") or "exec_failed")}
|
||
return {"ok": True, "data": data}
|
||
|
||
return ToolSpec(
|
||
name="netx_exec_managed_ne",
|
||
description=(
|
||
"经 netx 登录「网元管理」中的设备并执行只读 CLI(POST /v1/managed-ne/exec)。"
|
||
"每条命令须以 show / display / ping / ping6 开头(只读查询与连通探测,禁止 get/traceroute/改配置等);"
|
||
"禁止管道符、分号及改配置类命令;单次最多 5 条;默认读超时 60s。"
|
||
"返回合并输出(含命令回显);失败时含 error/detail。"
|
||
"先 netx_list_managed_ne 解析 ne_id;若 connect_status 非 pass 可先 netx_get_managed_ne 看 connect_detail。"
|
||
),
|
||
parameters={
|
||
"type": "object",
|
||
"properties": {
|
||
"ne_id": {"type": "string", "description": "纳管网元 UUID"},
|
||
"commands": {
|
||
"type": "array",
|
||
"items": {"type": "string"},
|
||
"minItems": 1,
|
||
"maxItems": 5,
|
||
"description": "只读 CLI 列表,如 show version、display interface brief",
|
||
},
|
||
"read_timeout_sec": {
|
||
"type": "integer",
|
||
"minimum": 10,
|
||
"maximum": 120,
|
||
"description": "单条命令 Netmiko 读超时(秒),默认 60",
|
||
},
|
||
},
|
||
"required": ["ne_id", "commands"],
|
||
"additionalProperties": False,
|
||
},
|
||
handler=handler,
|
||
tags=frozenset({"netx", "ops", "managed_ne", "cli", "exec"}),
|
||
risk_level="medium",
|
||
read_only=False,
|
||
)
|
||
|
||
|
||
__all__: list[str] = []
|
||
|
||
_LEGACY_EXPORTS = [
|
||
"netx_query_ume_alarms_tool",
|
||
"netx_aggregate_ume_alarms_tool",
|
||
"netx_run_ume_diagnostics_tool",
|
||
"netx_query_ume_ne_inventory_tool",
|
||
"netx_get_ume_ne_tool",
|
||
"netx_query_ume_alarms_raw_tool",
|
||
"netx_aggregate_ume_alarms_raw_tool",
|
||
"netx_list_ume_alarm_fields_tool",
|
||
"netx_sql_query_ume_tool",
|
||
"netx_list_managed_ne_tool",
|
||
"netx_get_managed_ne_tool",
|
||
"netx_exec_managed_ne_tool",
|
||
]
|
||
|
||
|
||
def _apply_legacy_exports() -> None:
|
||
import os as _os
|
||
|
||
if str(_os.getenv("OCLAW_NETX_BUILTIN_TOOLS") or "0").strip().lower() in {"0", "false", "no", "off"}:
|
||
return
|
||
globals()["__all__"] = list(_LEGACY_EXPORTS)
|
||
|
||
|
||
_apply_legacy_exports()
|