From 45ce9d3f7e001bc95e6ccc0d5eb967e9d996d1c9 Mon Sep 17 00:00:00 2001 From: oliver Date: Sat, 30 May 2026 15:44:32 +0800 Subject: [PATCH] feat(mcp): add netx-mcp package and HTTP stdio MCP - Extract installable packages/netx-mcp (12 HTTP tools, mcp.json) - Delegate netx_api.mcp to netx_mcp; keep legacy db_server shim - Add docs/MCP.md, install payload, pyproject entry, tests Co-authored-by: Cursor --- .gitignore | 1 + README.md | 18 +- docs/MCP.md | 219 +++++++ mcp.json | 15 + mcp_install_payload.json | 29 +- netx_api/__init__.py | 4 +- netx_api/mcp/__init__.py | 5 + netx_api/mcp/__main__.py | 4 + netx_api/mcp/db_server.py | 561 ++++++++++++++++++ netx_api/mcp/http_client.py | 3 + netx_api/mcp/http_tools.py | 3 + netx_api/mcp_server.py | 560 +---------------- packages/netx-mcp/README.md | 47 ++ packages/netx-mcp/mcp.json | 15 + packages/netx-mcp/pyproject.toml | 20 + packages/netx-mcp/src/netx_mcp/__init__.py | 5 + packages/netx-mcp/src/netx_mcp/__main__.py | 4 + packages/netx-mcp/src/netx_mcp/http_client.py | 107 ++++ packages/netx-mcp/src/netx_mcp/http_tools.py | 428 +++++++++++++ packages/netx-mcp/src/netx_mcp/server.py | 87 +++ packages/netx-mcp/tests/test_mcp_http.py | 87 +++ pyproject.toml | 37 ++ tests/test_mcp_http.py | 107 ++++ 23 files changed, 1810 insertions(+), 556 deletions(-) create mode 100644 docs/MCP.md create mode 100644 mcp.json create mode 100644 netx_api/mcp/__init__.py create mode 100644 netx_api/mcp/__main__.py create mode 100644 netx_api/mcp/db_server.py create mode 100644 netx_api/mcp/http_client.py create mode 100644 netx_api/mcp/http_tools.py create mode 100644 packages/netx-mcp/README.md create mode 100644 packages/netx-mcp/mcp.json create mode 100644 packages/netx-mcp/pyproject.toml create mode 100644 packages/netx-mcp/src/netx_mcp/__init__.py create mode 100644 packages/netx-mcp/src/netx_mcp/__main__.py create mode 100644 packages/netx-mcp/src/netx_mcp/http_client.py create mode 100644 packages/netx-mcp/src/netx_mcp/http_tools.py create mode 100644 packages/netx-mcp/src/netx_mcp/server.py create mode 100644 packages/netx-mcp/tests/test_mcp_http.py create mode 100644 pyproject.toml create mode 100644 tests/test_mcp_http.py diff --git a/.gitignore b/.gitignore index 539a603..c2f6b4d 100644 --- a/.gitignore +++ b/.gitignore @@ -2,6 +2,7 @@ __pycache__/ *.pyc .pytest_cache/ +*.egg-info/ .env build/ dist/ diff --git a/README.md b/README.md index e670991..bd4c91b 100644 --- a/README.md +++ b/README.md @@ -125,15 +125,23 @@ powershell -ExecutionPolicy Bypass -File .\scripts\stop_netx.ps1 -Force Primary web UI (Vite): `http://127.0.0.1:5173/` API base: `http://127.0.0.1:8890/` -### 6) Optional: MCP server (for oclaw integration) +### 6) MCP(Cursor / oclaw / Claude) + +先启动 netx API(§5),再在 **MCP 宿主同机** 安装轻量客户端并配置。 + +**完整说明(安装、配置、更新、排错)见:[docs/MCP.md](docs/MCP.md)** + +速查: ```powershell -.\.venv\Scripts\python -m netx_api.mcp_server +pip install -e ./packages/netx-mcp +# 配置见 mcp.json,运行: +python -m netx_mcp ``` -Or install via payload: - -- `netx/mcp_install_payload.json` +- 客户端配置:[`mcp.json`](mcp.json)(Cursor / oclaw Admin 粘贴同一份) +- oclaw 可选 payload:[`mcp_install_payload.json`](mcp_install_payload.json) +- 子包说明:[`packages/netx-mcp/README.md`](packages/netx-mcp/README.md) ## Useful API endpoints diff --git a/docs/MCP.md b/docs/MCP.md new file mode 100644 index 0000000..6159cc0 --- /dev/null +++ b/docs/MCP.md @@ -0,0 +1,219 @@ +# netx MCP — 安装与更新 + +netx 通过 **stdio MCP** 把告警/网元能力暴露给 Cursor、Claude Desktop、oclaw 等宿主。MCP 进程只做 HTTP 客户端,**不直连数据库**;数据来自已运行的 netx REST API。 + +``` +MCP 宿主 (Cursor / oclaw) → stdio netx_mcp → HTTP NETX_API_URL → netx API +``` + +| 组件 | 部署位置 | 说明 | +|------|----------|------| +| **netx API** | 本机或远端服务器 | `python -m netx_api.main`,默认 `http://127.0.0.1:8890` | +| **netx-mcp** | 与 MCP 宿主同机 | `pip install` 轻量包,仅需 Python + httpx | + +--- + +## 前置条件 + +1. **netx API 已启动**且可访问: + + ```powershell + curl http://127.0.0.1:8890/health + ``` + +2. **Python 3.11+**(与 MCP 宿主使用的 `python` 一致)。 + +3. 远端 API 时:记下 base URL(无尾部 `/`),例如 `http://10.0.0.5:8890`。 + +--- + +## 1. 安装 `netx-mcp` 包 + +在 **运行 MCP 的那台机器**上执行(不必安装整套 netx 服务依赖): + +```powershell +cd D:\project\chatgpt\netx +pip install -e ./packages/netx-mcp +``` + +验证: + +```powershell +python -m netx_mcp +# 另开终端发 JSON-RPC initialize(或见下文「自检」) +python -c "import netx_mcp; print('ok')" +``` + +从 GitHub 安装(无本地仓库时,**子目录必须是 `packages/netx-mcp`**): + +```powershell +pip install "git+https://github.com/hansjone/netx.git#subdirectory=packages/netx-mcp" +``` + +固定分支或 tag 时在 `.git` 后加 `@`,例如 `@main` 或 `@v0.2.0`: + +```powershell +pip install "git+https://github.com/hansjone/netx.git@main#subdirectory=packages/netx-mcp" +``` + +私有仓库需本机已配置 Git 凭据(或 SSH:`git+ssh://git@github.com/hansjone/netx.git#subdirectory=packages/netx-mcp`)。 + +开发者在 netx 仓库根目录也可 `pip install -e .`(含 API);MCP 仍推荐只装 `packages/netx-mcp`。 + +--- + +## 2. 环境变量 + +| 变量 | 必填 | 默认 | 说明 | +|------|------|------|------| +| `NETX_API_URL` | 否 | `http://127.0.0.1:8890` | netx REST 根地址,可指向远端 | +| `NETX_API_TOKEN` | 否 | 空 | API 启用 Bearer 时填写 | +| `NETX_LANG` | 否 | `zh` | `zh` / `en`,影响 API 文案 | + +本机默认端口时 **可不设任何变量**。 + +--- + +## 3. 在 Cursor / Claude Desktop 中配置 + +复制仓库中的 [`mcp.json`](../mcp.json)(或 [`packages/netx-mcp/mcp.json`](../packages/netx-mcp/mcp.json))到客户端 MCP 配置,例如 Cursor:`.cursor/mcp.json`。 + +```json +{ + "mcpServers": { + "netx": { + "command": "python", + "args": ["-m", "netx_mcp"], + "env": { + "NETX_API_URL": "http://127.0.0.1:8890", + "NETX_API_TOKEN": "", + "NETX_LANG": "zh", + "PYTHONIOENCODING": "utf-8", + "PYTHONUTF8": "1" + } + } + } +} +``` + +保存后 **重启 Cursor/客户端**,使 MCP 子进程重新拉起。 + +--- + +## 4. 在 oclaw Admin 中安装 + +与 Cursor **同一份** `mcpServers` JSON 即可,无需转成别的格式。 + +1. 完成上文 **§1**(本机 `pip install -e packages/netx-mcp`)。 +2. Admin → MCP → **安装 JSON**,粘贴 `mcp.json` 全文。 +3. **Health** → **Sync Tools**(应看到 **12** 个工具)。 +4. 在 **MCP 专家绑定** 中为 ops 专家勾选 `server_id=netx`。 + +更细的 oclaw 说明(双轨内置工具、锚点注入等)见 oclaw 仓库: +`oclaw/docs/NETX_MCP_INTEGRATION.md`。 + +可选:oclaw 专用字段展开版 [`mcp_install_payload.json`](../mcp_install_payload.json)(与 `mcp.json` 等价)。 + +--- + +## 5. 暴露的 12 个工具 + +| 类别 | 工具名 | +|------|--------| +| UME 告警 | `queryUmeAlarms`, `aggregateUmeAlarms`, `runUmeDiagnostics` | +| UME 网元 | `queryUmeNeInventory`, `getUmeNe` | +| UME 原始/SQL | `queryUmeAlarmsRaw`, `aggregateUmeAlarmsRaw`, `listUmeAlarmFields`, `sqlQueryUme` | +| 托管网元 CLI | `listManagedNe`, `getManagedNe`, `execManagedNe` | + +oclaw 中名称带前缀:`mcp__netx__`。 + +--- + +## 6. 更新 MCP + +MCP 代码在 **`packages/netx-mcp`**(版本见 `packages/netx-mcp/pyproject.toml`)。更新步骤: + +### 6.1 拉代码并重装包 + +```powershell +cd D:\project\chatgpt\netx +git pull +pip install -e ./packages/netx-mcp +``` + +确认版本(可选): + +```powershell +pip show netx-mcp +``` + +### 6.2 重启 MCP 宿主 + +| 宿主 | 操作 | +|------|------| +| **Cursor / Claude** | 完全退出客户端后重开,或重载 MCP | +| **oclaw** | 重启 gateway/主进程;Admin 对 `netx` 再点 **Health** → **Sync Tools** | + +仅改 `NETX_API_URL` 等 env 时,同样需重启 MCP 子进程(或 oclaw 整进程)。 + +### 6.3 无需改配置的情况 + +- 只改 **netx API** 业务逻辑、未改 MCP 工具名/参数:重装 `netx-mcp` 可选;API 部署后 MCP 自动走新 API。 +- 改了 **工具列表或 JSON schema**:必须重装 `netx-mcp` 并在宿主 **Sync Tools**。 + +### 6.4 从旧入口迁移 + +| 旧方式 | 新方式 | +|--------|--------| +| `python -m netx_api.mcp` | `python -m netx_mcp`(推荐) | +| `NETX_MCP_MODE=db` 直连库 | 已废弃;使用 HTTP 模式 | + +根目录 `python -m netx_api.mcp` 仍会委托到 `netx_mcp`,新环境请只装 `netx-mcp` 包。 + +--- + +## 7. 自检 + +**API:** + +```powershell +curl http://127.0.0.1:8890/health +``` + +**MCP 工具列表(需已 `pip install -e packages/netx-mcp`):** + +```powershell +cd D:\project\chatgpt\netx +python -m pytest packages/netx-mcp/tests/test_mcp_http.py -q +``` + +**手动 stdio(PowerShell 示例):** + +```powershell +$env:NETX_API_URL = "http://127.0.0.1:8890" +$p = Start-Process python -ArgumentList "-m","netx_mcp" -RedirectStandardInput pipe -RedirectStandardOutput pipe -NoNewWindow -PassThru +# 向 stdin 写入一行 JSON-RPC initialize / tools/list(见 packages/netx-mcp/tests) +``` + +--- + +## 8. 常见问题 + +| 现象 | 处理 | +|------|------| +| `ModuleNotFoundError: netx_mcp` | 在 MCP 宿主使用的 Python 上执行 `pip install -e packages/netx-mcp` | +| 工具调用连不上 API | 检查 `NETX_API_URL`、防火墙、远端 API 是否启动 | +| Windows 乱码 / JSON 解析失败 | 配置里保留 `PYTHONIOENCODING=utf-8`、`PYTHONUTF8=1`(见 `mcp.json`) | +| oclaw 工具数为 0 | Admin **Sync Tools**;确认 `server_id=netx` 已绑定专家 | +| 仍想用仓库脚本路径 | 未 pip 安装时可临时 `"args": ["D:/.../netx_api/mcp_server.py"]`(开发用) | + +--- + +## 相关文件 + +| 文件 | 用途 | +|------|------| +| [`packages/netx-mcp/`](../packages/netx-mcp/) | MCP 实现与 `pyproject.toml` | +| [`mcp.json`](../mcp.json) | Cursor / oclaw 粘贴用配置 | +| [`mcp_install_payload.json`](../mcp_install_payload.json) | oclaw 字段展开版(可选) | +| [`packages/netx-mcp/README.md`](../packages/netx-mcp/README.md) | 子包速查 | diff --git a/mcp.json b/mcp.json new file mode 100644 index 0000000..4ed18a7 --- /dev/null +++ b/mcp.json @@ -0,0 +1,15 @@ +{ + "mcpServers": { + "netx": { + "command": "python", + "args": ["-m", "netx_mcp"], + "env": { + "NETX_API_URL": "http://127.0.0.1:8890", + "NETX_API_TOKEN": "", + "NETX_LANG": "zh", + "PYTHONIOENCODING": "utf-8", + "PYTHONUTF8": "1" + } + } + } +} diff --git a/mcp_install_payload.json b/mcp_install_payload.json index 32e71ce..88bd29c 100644 --- a/mcp_install_payload.json +++ b/mcp_install_payload.json @@ -1,15 +1,28 @@ { "source_type": "local", - "source_ref": "netx-local-mcp", - "server_id": "netx-local", - "version": "0.1.0", + "source_ref": "netx-mcp", + "server_id": "netx", + "version": "0.2.0", "entry_command": "python", - "entry_args": [ - "D:/project/chatgpt/netx/netx_api/mcp_server.py" - ], - "env_schema": {}, + "entry_args": ["-m", "netx_mcp"], + "env_schema": { + "NETX_API_URL": { + "type": "string", + "description": "netx REST API base URL (no trailing slash)", + "default": "http://127.0.0.1:8890" + }, + "NETX_API_TOKEN": { + "type": "string", + "description": "Optional Bearer token when netx API auth is enabled" + }, + "NETX_LANG": { + "type": "string", + "description": "Response language hint: zh or en", + "default": "zh" + } + }, "required_permissions": [], "risk_level": "medium", "enabled": true, - "timeout_s": 30 + "timeout_s": 120 } diff --git a/netx_api/__init__.py b/netx_api/__init__.py index a05eb9a..31e8971 100644 --- a/netx_api/__init__.py +++ b/netx_api/__init__.py @@ -1,3 +1 @@ -__all__ = ["__version__"] - -__version__ = "0.1.0" +"""Run netx API or MCP entrypoints via ``python -m netx_api.``.""" diff --git a/netx_api/mcp/__init__.py b/netx_api/mcp/__init__.py new file mode 100644 index 0000000..35f01fb --- /dev/null +++ b/netx_api/mcp/__init__.py @@ -0,0 +1,5 @@ +"""Backward-compatible re-exports; prefer ``pip install netx-mcp`` and ``import netx_mcp``.""" + +from netx_mcp.http_tools import HTTP_MCP_TOOLS, call_http_tool + +__all__ = ["HTTP_MCP_TOOLS", "call_http_tool"] diff --git a/netx_api/mcp/__main__.py b/netx_api/mcp/__main__.py new file mode 100644 index 0000000..f3a2a9b --- /dev/null +++ b/netx_api/mcp/__main__.py @@ -0,0 +1,4 @@ +from netx_mcp.server import main + +if __name__ == "__main__": + main() diff --git a/netx_api/mcp/db_server.py b/netx_api/mcp/db_server.py new file mode 100644 index 0000000..71bf228 --- /dev/null +++ b/netx_api/mcp/db_server.py @@ -0,0 +1,561 @@ +from __future__ import annotations + +import json +import sys +from typing import Any + +from .db import Base, SessionLocal, engine +from .importer import aggregate_alarms, query_alarms +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: + sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": rid, "result": result}, ensure_ascii=False) + "\n") + sys.stdout.flush() + + +def _err(rid: Any, code: int, message: str) -> None: + sys.stdout.write( + json.dumps({"jsonrpc": "2.0", "id": rid, "error": {"code": code, "message": message}}, ensure_ascii=False) + + "\n" + ) + sys.stdout.flush() + + +def _tool_list() -> list[dict[str, Any]]: + return [ + { + "name": "queryAlarms", + "description": "Query normalized alarms with filters and pagination.", + "inputSchema": { + "type": "object", + "properties": { + "batch_id": {"type": "string"}, + "alarm_code": {"type": "string"}, + "severity": {"type": "string"}, + "ne_name": {"type": "string"}, + "page": {"type": "integer", "minimum": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 200}, + }, + "additionalProperties": False, + }, + }, + { + "name": "aggregateAlarms", + "description": "Aggregate alarms by severity_norm, alarm_code or ne_name.", + "inputSchema": { + "type": "object", + "properties": { + "group_by": { + "type": "string", + "enum": ["severity_norm", "alarm_code", "ne_name"], + }, + "batch_id": {"type": "string"}, + }, + "required": ["group_by"], + "additionalProperties": False, + }, + }, + { + "name": "getImportBatch", + "description": "Get imported batch summary by batch_id.", + "inputSchema": { + "type": "object", + "properties": {"batch_id": {"type": "string"}}, + "required": ["batch_id"], + "additionalProperties": False, + }, + }, + { + "name": "runDiagnostics", + "description": "Generate quick diagnostics summary by batch.", + "inputSchema": { + "type": "object", + "properties": {"batch_id": {"type": "string"}}, + "required": ["batch_id"], + "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, + }, + }, + ] + + +def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: + with SessionLocal() as db: + if name == "queryAlarms": + total, rows = query_alarms( + db, + batch_id=str(args.get("batch_id") or "").strip() or None, + alarm_code=str(args.get("alarm_code") or "").strip() or None, + severity=str(args.get("severity") or "").strip() or None, + ne_name=str(args.get("ne_name") or "").strip() or None, + page=max(1, int(args.get("page") or 1)), + page_size=min(200, max(1, int(args.get("page_size") or 50))), + ) + payload = { + "total": total, + "items": [ + { + "id": r.id, + "batch_id": r.batch_id, + "alarm_time": r.alarm_time.isoformat(), + "severity_norm": r.severity_norm, + "ne_name": r.ne_name, + "alarm_code": r.alarm_code, + "alarm_name": r.alarm_name, + } + for r in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "aggregateAlarms": + group_by = str(args.get("group_by") or "").strip() + rows = aggregate_alarms( + db, + group_by=group_by, + batch_id=str(args.get("batch_id") or "").strip() or None, + ) + payload = {"group_by": group_by, "buckets": [{"key": k, "count": v} for k, v in rows]} + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "getImportBatch": + batch_id = str(args.get("batch_id") or "").strip() + batch = db.get(AlarmBatch, batch_id) + if not batch: + return {"content": [{"type": "text", "text": json.dumps({"error": "batch_not_found"})}], "isError": True} + payload = { + "batch_id": batch.batch_id, + "total_rows": batch.total_rows, + "success_rows": batch.success_rows, + "failed_rows": batch.failed_rows, + "status": batch.status, + "created_at": batch.created_at.isoformat(), + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "runDiagnostics": + batch_id = str(args.get("batch_id") or "").strip() + sev_rows = aggregate_alarms(db, group_by="severity_norm", batch_id=batch_id) + code_rows = aggregate_alarms(db, group_by="alarm_code", batch_id=batch_id)[:5] + ne_rows = aggregate_alarms(db, group_by="ne_name", batch_id=batch_id)[:5] + sev_map = {k: v for k, v in sev_rows} + findings: list[str] = [] + actions: list[str] = [] + risk_level = "low" + if int(sev_map.get("critical", 0)) > 0: + findings.append("critical 告警存在,建议优先确认核心网元影响。") + actions.append("优先处理 critical 告警,确认影响面并升级。") + risk_level = "high" + if int(sev_map.get("warning", 0)) > int(sev_map.get("major", 0)): + findings.append("warning 占比较高,疑似阈值型告警风暴。") + actions.append("检查高频 warning 告警码是否集中在单一阈值策略。") + if risk_level != "high": + risk_level = "medium" + if not findings: + findings.append("分布相对均衡,建议按 top 告警码进一步排查。") + actions.append("按 top 告警码和网元继续下钻分析。") + payload = { + "batch_id": batch_id, + "risk_level": risk_level, + "severity_summary": [{"key": k, "count": v} for k, v in sev_rows], + "top_alarm_codes": [{"key": k, "count": v} for k, v in code_rows], + "top_ne": [{"key": k, "count": v} for k, v in ne_rows], + "findings": findings, + "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, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id + ) + 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.host_name.contains(keyword) + | UmeInventoryNE.ne_name.contains(keyword) + | UmeInventoryNE.user_label.contains(keyword) + | UmeInventoryNE.ip_address.contains(keyword) + | UmeInventoryNE.host_name.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.time_created.desc(), + UmeAlarmCurrent.last_seen_at.desc(), + UmeAlarmCurrent.alarm_key.desc(), + ) + .offset((page - 1) * page_size) + .limit(page_size) + .all() + ) + payload = { + "total": total, + "page": page, + "page_size": page_size, + "items": [ + { + "alarm_key": str(alarm.alarm_key or ""), + "ne_id": str(alarm.ne_id or ""), + "host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), + "object_name": str(alarm.object_name or ""), + "event_type": str(alarm.event_type or ""), + "native_probable_cause": str(alarm.native_probable_cause or ""), + "perceived_severity": str(alarm.perceived_severity or ""), + "is_cleared": str(alarm.is_cleared or ""), + "time_created": str(alarm.time_created or ""), + } + for alarm, ne in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if name == "umeListHistoryAlarms": + stmt = db.query(UmeAlarmHistory, UmeInventoryNE).outerjoin( + UmeInventoryNE, UmeAlarmHistory.ne_id == UmeInventoryNE.ne_id + ) + 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.host_name.contains(keyword) + | UmeInventoryNE.ne_name.contains(keyword) + | UmeInventoryNE.user_label.contains(keyword) + | UmeInventoryNE.ip_address.contains(keyword) + | UmeInventoryNE.host_name.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(alarm.alarm_key or ""), + "ne_id": str(alarm.ne_id or ""), + "host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""), + "ne_name": str((ne.ne_name if ne else "") or ""), + "user_label": str((ne.user_label if ne else "") or ""), + "object_name": str(alarm.object_name or ""), + "event_type": str(alarm.event_type or ""), + "native_probable_cause": str(alarm.native_probable_cause or ""), + "perceived_severity": str(alarm.perceived_severity or ""), + "is_cleared": str(alarm.is_cleared or ""), + "time_created": str(alarm.time_created or ""), + } + for alarm, ne in rows + ], + } + return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + raise ValueError(f"unknown tool: {name}") + + +def _ensure_db_schema() -> None: + try: + Base.metadata.create_all(bind=engine) + with engine.begin() as conn: + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS device_level VARCHAR(64) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS location VARCHAR(512) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS ipv6_address VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql( + "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS hardware_version VARCHAR(128) DEFAULT ''" + ) + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS loopback VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql( + "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS consistent_state VARCHAR(64) DEFAULT ''" + ) + conn.exec_driver_sql( + "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS interface_version VARCHAR(128) DEFAULT ''" + ) + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS mac VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS admin_status VARCHAR(64) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS address_type VARCHAR(64) DEFAULT ''") + conn.exec_driver_sql( + "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS connection_status VARCHAR(64) DEFAULT ''" + ) + conn.exec_driver_sql( + "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS maintain_status VARCHAR(64) DEFAULT ''" + ) + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS net_mask VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS create_time VARCHAR(64) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS creator VARCHAR(128) DEFAULT ''") + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS user_label") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS ne_name") + conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS user_label") + conn.exec_driver_sql("COMMENT ON TABLE ume_inventory_ne IS '网元对象详细信息'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_id IS '网元uuid'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_name IS '资源名称'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_type IS '网元类型'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.user_label IS '用户标签'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.address_type IS '管理地址类型(1:IPv4,2:IPv6)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ip_address IS '网元IPv4地址'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.net_mask IS '管理IPv4掩码(点分十进制)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ipv6_address IS 'IPv6地址'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.admin_status IS '管理状态(0-离线,1-在线)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.connection_status IS '连接状态(0-断链,1-正常)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.consistent_state IS '数据一致性状态(1一致,2不一致,3冲突)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.maintain_status IS '工程状态(0普通,1调测,2新建)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.vendor IS '网元提供商'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.interface_version IS '网元接口版本号'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.hardware_version IS '硬件版本'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.mac IS '设备机架MAC地址'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.loopback IS '业务环回IP(IPv4)'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.device_level IS '网元层次'") + conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.host_name IS '主机名称'") + except Exception: + pass + + +def run_stdio_loop() -> None: + """Legacy MCP: direct PostgreSQL access (deprecated; use HTTP mode).""" + _ensure_db_schema() + for line in sys.stdin: + raw = line.strip() + if not raw: + continue + try: + req = json.loads(raw) + except Exception: + continue + rid = req.get("id") + method = str(req.get("method") or "") + params = req.get("params") if isinstance(req.get("params"), dict) else {} + + try: + if method == "initialize": + _ok( + rid, + { + "protocolVersion": "2024-11-05", + "capabilities": {"tools": {}}, + "serverInfo": {"name": "netx-mcp", "version": "0.1.0"}, + }, + ) + continue + if method == "notifications/initialized": + continue + if method == "tools/list": + _ok(rid, {"tools": _tool_list()}) + continue + if method == "tools/call": + name = str(params.get("name") or "") + args = params.get("arguments") if isinstance(params.get("arguments"), dict) else {} + _ok(rid, _call_tool(name, args)) + continue + _err(rid, -32601, f"method not found: {method}") + except Exception as exc: + _err(rid, -32000, str(exc)) + + +def main() -> None: + run_stdio_loop() + + +if __name__ == "__main__": + main() diff --git a/netx_api/mcp/http_client.py b/netx_api/mcp/http_client.py new file mode 100644 index 0000000..0288d66 --- /dev/null +++ b/netx_api/mcp/http_client.py @@ -0,0 +1,3 @@ +"""Re-export from ``netx-mcp`` package (install: ``pip install -e packages/netx-mcp``).""" + +from netx_mcp.http_client import * # noqa: F403 diff --git a/netx_api/mcp/http_tools.py b/netx_api/mcp/http_tools.py new file mode 100644 index 0000000..d77966c --- /dev/null +++ b/netx_api/mcp/http_tools.py @@ -0,0 +1,3 @@ +"""Re-export from ``netx-mcp`` package (install: ``pip install -e packages/netx-mcp``).""" + +from netx_mcp.http_tools import * # noqa: F403 diff --git a/netx_api/mcp_server.py b/netx_api/mcp_server.py index e6690de..4f31601 100644 --- a/netx_api/mcp_server.py +++ b/netx_api/mcp_server.py @@ -1,551 +1,31 @@ +"""Backward-compatible entry for ``python netx_api/mcp_server.py`` / ``netx-mcp`` script on netx-ops install. + +HTTP MCP: delegates to ``netx_mcp`` (``pip install -e packages/netx-mcp``). +Legacy DB mode (``NETX_MCP_MODE=db``): requires full ``netx-ops`` install only. +""" + from __future__ import annotations -import json +import os import sys -from typing import Any +from pathlib import Path -from .db import Base, SessionLocal, engine -from .importer import aggregate_alarms, query_alarms -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: - sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": rid, "result": result}, ensure_ascii=False) + "\n") - sys.stdout.flush() - - -def _err(rid: Any, code: int, message: str) -> None: - sys.stdout.write( - json.dumps({"jsonrpc": "2.0", "id": rid, "error": {"code": code, "message": message}}, ensure_ascii=False) - + "\n" - ) - sys.stdout.flush() - - -def _tool_list() -> list[dict[str, Any]]: - return [ - { - "name": "queryAlarms", - "description": "Query normalized alarms with filters and pagination.", - "inputSchema": { - "type": "object", - "properties": { - "batch_id": {"type": "string"}, - "alarm_code": {"type": "string"}, - "severity": {"type": "string"}, - "ne_name": {"type": "string"}, - "page": {"type": "integer", "minimum": 1}, - "page_size": {"type": "integer", "minimum": 1, "maximum": 200}, - }, - "additionalProperties": False, - }, - }, - { - "name": "aggregateAlarms", - "description": "Aggregate alarms by severity_norm, alarm_code or ne_name.", - "inputSchema": { - "type": "object", - "properties": { - "group_by": { - "type": "string", - "enum": ["severity_norm", "alarm_code", "ne_name"], - }, - "batch_id": {"type": "string"}, - }, - "required": ["group_by"], - "additionalProperties": False, - }, - }, - { - "name": "getImportBatch", - "description": "Get imported batch summary by batch_id.", - "inputSchema": { - "type": "object", - "properties": {"batch_id": {"type": "string"}}, - "required": ["batch_id"], - "additionalProperties": False, - }, - }, - { - "name": "runDiagnostics", - "description": "Generate quick diagnostics summary by batch.", - "inputSchema": { - "type": "object", - "properties": {"batch_id": {"type": "string"}}, - "required": ["batch_id"], - "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, - }, - }, - ] - - -def _call_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: - with SessionLocal() as db: - if name == "queryAlarms": - total, rows = query_alarms( - db, - batch_id=str(args.get("batch_id") or "").strip() or None, - alarm_code=str(args.get("alarm_code") or "").strip() or None, - severity=str(args.get("severity") or "").strip() or None, - ne_name=str(args.get("ne_name") or "").strip() or None, - page=max(1, int(args.get("page") or 1)), - page_size=min(200, max(1, int(args.get("page_size") or 50))), - ) - payload = { - "total": total, - "items": [ - { - "id": r.id, - "batch_id": r.batch_id, - "alarm_time": r.alarm_time.isoformat(), - "severity_norm": r.severity_norm, - "ne_name": r.ne_name, - "alarm_code": r.alarm_code, - "alarm_name": r.alarm_name, - } - for r in rows - ], - } - return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} - if name == "aggregateAlarms": - group_by = str(args.get("group_by") or "").strip() - rows = aggregate_alarms( - db, - group_by=group_by, - batch_id=str(args.get("batch_id") or "").strip() or None, - ) - payload = {"group_by": group_by, "buckets": [{"key": k, "count": v} for k, v in rows]} - return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} - if name == "getImportBatch": - batch_id = str(args.get("batch_id") or "").strip() - batch = db.get(AlarmBatch, batch_id) - if not batch: - return {"content": [{"type": "text", "text": json.dumps({"error": "batch_not_found"})}], "isError": True} - payload = { - "batch_id": batch.batch_id, - "total_rows": batch.total_rows, - "success_rows": batch.success_rows, - "failed_rows": batch.failed_rows, - "status": batch.status, - "created_at": batch.created_at.isoformat(), - } - return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} - if name == "runDiagnostics": - batch_id = str(args.get("batch_id") or "").strip() - sev_rows = aggregate_alarms(db, group_by="severity_norm", batch_id=batch_id) - code_rows = aggregate_alarms(db, group_by="alarm_code", batch_id=batch_id)[:5] - ne_rows = aggregate_alarms(db, group_by="ne_name", batch_id=batch_id)[:5] - sev_map = {k: v for k, v in sev_rows} - findings: list[str] = [] - actions: list[str] = [] - risk_level = "low" - if int(sev_map.get("critical", 0)) > 0: - findings.append("critical 告警存在,建议优先确认核心网元影响。") - actions.append("优先处理 critical 告警,确认影响面并升级。") - risk_level = "high" - if int(sev_map.get("warning", 0)) > int(sev_map.get("major", 0)): - findings.append("warning 占比较高,疑似阈值型告警风暴。") - actions.append("检查高频 warning 告警码是否集中在单一阈值策略。") - if risk_level != "high": - risk_level = "medium" - if not findings: - findings.append("分布相对均衡,建议按 top 告警码进一步排查。") - actions.append("按 top 告警码和网元继续下钻分析。") - payload = { - "batch_id": batch_id, - "risk_level": risk_level, - "severity_summary": [{"key": k, "count": v} for k, v in sev_rows], - "top_alarm_codes": [{"key": k, "count": v} for k, v in code_rows], - "top_ne": [{"key": k, "count": v} for k, v in ne_rows], - "findings": findings, - "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, UmeInventoryNE).outerjoin( - UmeInventoryNE, UmeAlarmCurrent.ne_id == UmeInventoryNE.ne_id - ) - 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.host_name.contains(keyword) - | UmeInventoryNE.ne_name.contains(keyword) - | UmeInventoryNE.user_label.contains(keyword) - | UmeInventoryNE.ip_address.contains(keyword) - | UmeInventoryNE.host_name.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.time_created.desc(), - UmeAlarmCurrent.last_seen_at.desc(), - UmeAlarmCurrent.alarm_key.desc(), - ) - .offset((page - 1) * page_size) - .limit(page_size) - .all() - ) - payload = { - "total": total, - "page": page, - "page_size": page_size, - "items": [ - { - "alarm_key": str(alarm.alarm_key or ""), - "ne_id": str(alarm.ne_id or ""), - "host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""), - "ne_name": str((ne.ne_name if ne else "") or ""), - "user_label": str((ne.user_label if ne else "") or ""), - "object_name": str(alarm.object_name or ""), - "event_type": str(alarm.event_type or ""), - "native_probable_cause": str(alarm.native_probable_cause or ""), - "perceived_severity": str(alarm.perceived_severity or ""), - "is_cleared": str(alarm.is_cleared or ""), - "time_created": str(alarm.time_created or ""), - } - for alarm, ne in rows - ], - } - return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} - if name == "umeListHistoryAlarms": - stmt = db.query(UmeAlarmHistory, UmeInventoryNE).outerjoin( - UmeInventoryNE, UmeAlarmHistory.ne_id == UmeInventoryNE.ne_id - ) - 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.host_name.contains(keyword) - | UmeInventoryNE.ne_name.contains(keyword) - | UmeInventoryNE.user_label.contains(keyword) - | UmeInventoryNE.ip_address.contains(keyword) - | UmeInventoryNE.host_name.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(alarm.alarm_key or ""), - "ne_id": str(alarm.ne_id or ""), - "host_name": str(alarm.host_name or (ne.host_name if ne else "") or ""), - "ne_name": str((ne.ne_name if ne else "") or ""), - "user_label": str((ne.user_label if ne else "") or ""), - "object_name": str(alarm.object_name or ""), - "event_type": str(alarm.event_type or ""), - "native_probable_cause": str(alarm.native_probable_cause or ""), - "perceived_severity": str(alarm.perceived_severity or ""), - "is_cleared": str(alarm.is_cleared or ""), - "time_created": str(alarm.time_created or ""), - } - for alarm, ne in rows - ], - } - return {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} - raise ValueError(f"unknown tool: {name}") +# Script-path launch without pip install netx-mcp: add packages/netx-mcp/src +_repo_root = Path(__file__).resolve().parent.parent +_mcp_src = _repo_root / "packages" / "netx-mcp" / "src" +if _mcp_src.is_dir() and str(_mcp_src) not in sys.path: + sys.path.insert(0, str(_mcp_src)) def main() -> None: - try: - Base.metadata.create_all(bind=engine) - with engine.begin() as conn: - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS device_level VARCHAR(64) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS location VARCHAR(512) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS ipv6_address VARCHAR(128) DEFAULT ''") - conn.exec_driver_sql( - "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS hardware_version VARCHAR(128) DEFAULT ''" - ) - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS loopback VARCHAR(128) DEFAULT ''") - conn.exec_driver_sql( - "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS consistent_state VARCHAR(64) DEFAULT ''" - ) - conn.exec_driver_sql( - "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS interface_version VARCHAR(128) DEFAULT ''" - ) - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS mac VARCHAR(128) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS admin_status VARCHAR(64) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS address_type VARCHAR(64) DEFAULT ''") - conn.exec_driver_sql( - "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS connection_status VARCHAR(64) DEFAULT ''" - ) - conn.exec_driver_sql( - "ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS maintain_status VARCHAR(64) DEFAULT ''" - ) - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS net_mask VARCHAR(128) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS create_time VARCHAR(64) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_inventory_ne ADD COLUMN IF NOT EXISTS creator VARCHAR(128) DEFAULT ''") - conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS ne_name") - conn.exec_driver_sql("ALTER TABLE ume_alarms_current DROP COLUMN IF EXISTS user_label") - conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS ne_name") - conn.exec_driver_sql("ALTER TABLE ume_alarms_history DROP COLUMN IF EXISTS user_label") - conn.exec_driver_sql("COMMENT ON TABLE ume_inventory_ne IS '网元对象详细信息'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_id IS '网元uuid'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_name IS '资源名称'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ne_type IS '网元类型'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.user_label IS '用户标签'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.address_type IS '管理地址类型(1:IPv4,2:IPv6)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ip_address IS '网元IPv4地址'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.net_mask IS '管理IPv4掩码(点分十进制)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.ipv6_address IS 'IPv6地址'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.admin_status IS '管理状态(0-离线,1-在线)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.connection_status IS '连接状态(0-断链,1-正常)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.consistent_state IS '数据一致性状态(1一致,2不一致,3冲突)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.maintain_status IS '工程状态(0普通,1调测,2新建)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.vendor IS '网元提供商'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.interface_version IS '网元接口版本号'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.hardware_version IS '硬件版本'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.mac IS '设备机架MAC地址'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.loopback IS '业务环回IP(IPv4)'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.device_level IS '网元层次'") - conn.exec_driver_sql("COMMENT ON COLUMN ume_inventory_ne.host_name IS '主机名称'") - except Exception: - pass - for line in sys.stdin: - raw = line.strip() - if not raw: - continue - try: - req = json.loads(raw) - except Exception: - continue - rid = req.get("id") - method = str(req.get("method") or "") - params = req.get("params") if isinstance(req.get("params"), dict) else {} + if str(os.getenv("NETX_MCP_MODE") or "http").strip().lower() == "db": + from netx_api.mcp.db_server import run_stdio_loop - try: - if method == "initialize": - _ok( - rid, - { - "protocolVersion": "2024-11-05", - "capabilities": {"tools": {}}, - "serverInfo": {"name": "netx-mcp", "version": "0.1.0"}, - }, - ) - continue - if method == "notifications/initialized": - continue - if method == "tools/list": - _ok(rid, {"tools": _tool_list()}) - continue - if method == "tools/call": - name = str(params.get("name") or "") - args = params.get("arguments") if isinstance(params.get("arguments"), dict) else {} - _ok(rid, _call_tool(name, args)) - continue - _err(rid, -32601, f"method not found: {method}") - except Exception as exc: - _err(rid, -32000, str(exc)) + run_stdio_loop() + return + from netx_mcp.server import main as mcp_main + + mcp_main() if __name__ == "__main__": diff --git a/packages/netx-mcp/README.md b/packages/netx-mcp/README.md new file mode 100644 index 0000000..455039a --- /dev/null +++ b/packages/netx-mcp/README.md @@ -0,0 +1,47 @@ +# netx-mcp + +轻量 **stdio MCP** 客户端,通过 HTTP 调用 [netx](../../README.md) REST API。装在 Cursor、oclaw、Claude Desktop 等 **MCP 宿主所在机器**;`NETX_API_URL` 可指向本机或远端 API。 + +**安装、更新、各宿主配置、排错** → 仓库主文档 **[docs/MCP.md](../../docs/MCP.md)**(请优先阅读)。 + +## 速查:安装 + +```powershell +# 在 netx 仓库根目录 +pip install -e ./packages/netx-mcp +python -c "import netx_mcp; print('ok')" +``` + +要求 **Python 3.11+**,且 netx API 已运行(`GET /health`)。 + +## 速查:更新 + +```powershell +cd +git pull +pip install -e ./packages/netx-mcp +``` + +然后 **重启 MCP 宿主**(Cursor 重开;oclaw 重启后在 Admin 对 `netx` 执行 Health → Sync Tools)。 + +## 运行 + +```powershell +$env:NETX_API_URL = "http://127.0.0.1:8890" +python -m netx_mcp +# 或: netx-mcp +``` + +## 配置 + +[`mcp.json`](./mcp.json) — `command: python`,`args: ["-m", "netx_mcp"]`,`env` 见文件。 + +## 工具(12) + +UME:`queryUmeAlarms`, `aggregateUmeAlarms`, `runUmeDiagnostics`, `queryUmeNeInventory`, `getUmeNe`, `queryUmeAlarmsRaw`, `aggregateUmeAlarmsRaw`, `listUmeAlarmFields`, `sqlQueryUme` + +托管网元:`listManagedNe`, `getManagedNe`, `execManagedNe` + +## 兼容 + +全量 `netx-ops` 开发安装下 `python -m netx_api.mcp` 仍会转到本包;新环境请只装 **netx-mcp**。 diff --git a/packages/netx-mcp/mcp.json b/packages/netx-mcp/mcp.json new file mode 100644 index 0000000..4ed18a7 --- /dev/null +++ b/packages/netx-mcp/mcp.json @@ -0,0 +1,15 @@ +{ + "mcpServers": { + "netx": { + "command": "python", + "args": ["-m", "netx_mcp"], + "env": { + "NETX_API_URL": "http://127.0.0.1:8890", + "NETX_API_TOKEN": "", + "NETX_LANG": "zh", + "PYTHONIOENCODING": "utf-8", + "PYTHONUTF8": "1" + } + } + } +} diff --git a/packages/netx-mcp/pyproject.toml b/packages/netx-mcp/pyproject.toml new file mode 100644 index 0000000..722e871 --- /dev/null +++ b/packages/netx-mcp/pyproject.toml @@ -0,0 +1,20 @@ +[build-system] +requires = ["setuptools>=68", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "netx-mcp" +version = "0.2.0" +description = "stdio MCP server for netx REST API (UME alarms, managed NE CLI)" +readme = "README.md" +requires-python = ">=3.11" +license = { text = "MIT" } +dependencies = [ + "httpx>=0.27.0", +] + +[project.scripts] +netx-mcp = "netx_mcp.server:main" + +[tool.setuptools.packages.find] +where = ["src"] diff --git a/packages/netx-mcp/src/netx_mcp/__init__.py b/packages/netx-mcp/src/netx_mcp/__init__.py new file mode 100644 index 0000000..fd8fab0 --- /dev/null +++ b/packages/netx-mcp/src/netx_mcp/__init__.py @@ -0,0 +1,5 @@ +"""netx-mcp: stdio MCP → netx HTTP API.""" + +from .http_tools import HTTP_MCP_TOOLS, call_http_tool + +__all__ = ["HTTP_MCP_TOOLS", "call_http_tool"] diff --git a/packages/netx-mcp/src/netx_mcp/__main__.py b/packages/netx-mcp/src/netx_mcp/__main__.py new file mode 100644 index 0000000..f3a2a9b --- /dev/null +++ b/packages/netx-mcp/src/netx_mcp/__main__.py @@ -0,0 +1,4 @@ +from netx_mcp.server import main + +if __name__ == "__main__": + main() diff --git a/packages/netx-mcp/src/netx_mcp/http_client.py b/packages/netx-mcp/src/netx_mcp/http_client.py new file mode 100644 index 0000000..0d312ee --- /dev/null +++ b/packages/netx-mcp/src/netx_mcp/http_client.py @@ -0,0 +1,107 @@ +"""HTTP client for netx REST API (used by stdio MCP server).""" + +from __future__ import annotations + +import json +import os +from typing import Any +from urllib.parse import quote + +import httpx + +_PROTOCOL_KEY_ZH_TO_EN: dict[str, str] = { + "其他": "Other", + "时钟": "Clock", + "OTN/光": "OTN/Optical", + "电源": "Power", +} + + +def api_base_url() -> str: + raw = ( + os.getenv("NETX_API_URL") + or os.getenv("OCLAW_NETX_BASE_URL") + or "http://127.0.0.1:8890" + ) + return str(raw or "").strip().rstrip("/") + + +def api_headers() -> dict[str, str]: + h = {"accept": "application/json"} + tok = (os.getenv("NETX_API_TOKEN") or os.getenv("OCLAW_NETX_API_TOKEN") or "").strip() + if tok: + h["authorization"] = f"Bearer {tok}" + return h + + +def lang_query_params() -> dict[str, str]: + lang = str(os.getenv("NETX_LANG") or "zh").strip().lower() + if lang.startswith("en"): + return {"lang": "en"} + return {} + + +def localize_payload(data: dict[str, Any]) -> dict[str, Any]: + lang = str(os.getenv("NETX_LANG") or "zh").strip().lower() + if not lang.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_json(method: str, path: str, *, params: dict[str, Any] | None = None, timeout: float = 45.0) -> dict[str, Any]: + url = f"{api_base_url()}{path}" + merged: dict[str, Any] = dict(lang_query_params()) + if params: + merged.update(params) + try: + with httpx.Client(timeout=timeout, trust_env=False) as client: + resp = client.request(method, url, params=merged or None, headers=api_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_payload(data) + 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_post_json(path: str, body: dict[str, Any], *, timeout: float = 180.0) -> dict[str, Any]: + url = f"{api_base_url()}{path}" + try: + with httpx.Client(timeout=timeout, trust_env=False) as client: + resp = client.post(url, json=body, headers=api_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_payload(data) + 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 mcp_text_result(payload: Any, *, is_error: bool = False) -> dict[str, Any]: + out: dict[str, Any] = {"content": [{"type": "text", "text": json.dumps(payload, ensure_ascii=False)}]} + if is_error: + out["isError"] = True + return out + + +def mcp_from_handler_result(result: dict[str, Any]) -> dict[str, Any]: + if not result.get("ok"): + return mcp_text_result(result, is_error=True) + return mcp_text_result(result) + + +def quote_ne_id(ne_id: str) -> str: + return quote(str(ne_id or "").strip(), safe="") diff --git a/packages/netx-mcp/src/netx_mcp/http_tools.py b/packages/netx-mcp/src/netx_mcp/http_tools.py new file mode 100644 index 0000000..ed4d620 --- /dev/null +++ b/packages/netx-mcp/src/netx_mcp/http_tools.py @@ -0,0 +1,428 @@ +"""MCP tool schemas and HTTP-backed handlers (12 tools: UME + managed NE; no Excel import batch).""" + +from __future__ import annotations + +from typing import Any, Callable + +from .http_client import http_json, http_post_json, mcp_from_handler_result, quote_ne_id + +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", +] + +_UME_RAW_FIELD_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 _query_ume_alarms(args: dict[str, Any]) -> dict[str, Any]: + 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() + keyword = str(args.get("keyword") or "").strip() + ne_name = str(args.get("ne_name") 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) + + +def _aggregate_ume_alarms(args: dict[str, Any]) -> dict[str, Any]: + _ = args + return http_json("GET", "/v1/ume/alarms/aggregate", params=None) + + +def _run_ume_diagnostics(args: dict[str, Any]) -> dict[str, Any]: + _ = args + return http_json("GET", "/v1/ume/diagnostics", params=None) + + +def _query_ume_ne_inventory(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) + + +def _get_ume_ne(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/ume/inventory/ne/{quote_ne_id(ne_id)}", params=None) + + +def _query_ume_alarms_raw(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(_UME_RAW_FIELD_PRESETS.get(preset) or []) + if fields: + params["select_fields"] = ",".join(fields) + return http_json("GET", "/v1/ume/alarms/raw", params=params) + + +def _aggregate_ume_alarms_raw(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) + + +def _list_ume_alarm_fields(args: dict[str, Any]) -> dict[str, Any]: + _ = args + return http_json("GET", "/v1/ume/alarms/fields", params=None) + + +def _sql_query_ume(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"} + return http_post_json( + "/v1/sql/ume_query", + {"sql": sql, "limit": limit, "statement_timeout_ms": statement_timeout_ms}, + timeout=60.0, + ) + + +def _list_managed_ne(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) + + +def _get_managed_ne(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) + + +def _exec_managed_ne(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} + + +HTTP_MCP_TOOLS: list[dict[str, Any]] = [ + { + "name": "queryUmeAlarms", + "description": "Query UME current alarms (each row includes host_name). Supports severity/ne_id/keyword and pagination.", + "inputSchema": { + "type": "object", + "properties": { + "severity": {"type": "string"}, + "ne_id": {"type": "string"}, + "ne_name": {"type": "string", "description": "Legacy alias mapped to keyword"}, + "keyword": {"type": "string"}, + "page": {"type": "integer", "minimum": 1, "default": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50}, + }, + "required": [], + "additionalProperties": False, + }, + }, + { + "name": "aggregateUmeAlarms", + "description": "Aggregate UME current alarms (by_severity/by_ne).", + "inputSchema": {"type": "object", "properties": {}, "required": [], "additionalProperties": False}, + }, + { + "name": "runUmeDiagnostics", + "description": "UME alarm diagnostics summary (severity distribution, top codes/NEs, protocol buckets).", + "inputSchema": {"type": "object", "properties": {}, "required": [], "additionalProperties": False}, + }, + { + "name": "queryUmeNeInventory", + "description": "Paged UME NE inventory synced in netx (keyword matches ne_id/ne_name/user_label/ip/host_name).", + "inputSchema": { + "type": "object", + "properties": { + "keyword": {"type": "string"}, + "page": {"type": "integer", "minimum": 1, "default": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50}, + }, + "required": [], + "additionalProperties": False, + }, + }, + { + "name": "getUmeNe", + "description": "Get single UME NE detail by ne_id (UUID).", + "inputSchema": { + "type": "object", + "properties": {"ne_id": {"type": "string"}}, + "required": ["ne_id"], + "additionalProperties": False, + }, + }, + { + "name": "queryUmeAlarmsRaw", + "description": "Power query UME current alarms with full alarm_* + ne_* fields; optional field_preset or select_fields.", + "inputSchema": { + "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"}, + "time_to": {"type": "string"}, + "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"}}, + "field_preset": {"type": "string", "enum": ["brief", "evidence", "ne_debug"]}, + "page": {"type": "integer", "minimum": 1, "default": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50}, + }, + "required": [], + "additionalProperties": False, + }, + }, + { + "name": "aggregateUmeAlarmsRaw", + "description": "Dynamic aggregation on UME raw fields (group_by/group_by2); prefer alarm_host_name for NE grouping.", + "inputSchema": { + "type": "object", + "properties": { + "group_by": {"type": "string", "enum": UME_RAW_GROUP_FIELDS}, + "group_by2": {"type": "string", "enum": UME_RAW_GROUP_FIELDS}, + "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, + }, + }, + { + "name": "listUmeAlarmFields", + "description": "List available fields for UME raw alarm queries.", + "inputSchema": {"type": "object", "properties": {}, "required": [], "additionalProperties": False}, + }, + { + "name": "sqlQueryUme", + "description": "Read-only SELECT on UME tables (ume_alarms_current/ume_inventory_ne); server enforces limits.", + "inputSchema": { + "type": "object", + "properties": { + "sql": {"type": "string"}, + "limit": {"type": "integer", "minimum": 1, "maximum": 2000, "default": 200}, + "statement_timeout_ms": {"type": "integer", "minimum": 0, "maximum": 30000, "default": 0}, + }, + "required": ["sql"], + "additionalProperties": False, + }, + }, + { + "name": "listManagedNe", + "description": "List netx managed NEs (SSH/Telnet inventory); use before execManagedNe.", + "inputSchema": { + "type": "object", + "properties": { + "keyword": {"type": "string"}, + "vendor": {"type": "string"}, + "connect_status": {"type": "string", "enum": ["unknown", "testing", "pass", "fail"]}, + "page": {"type": "integer", "minimum": 1, "default": 1}, + "page_size": {"type": "integer", "minimum": 1, "maximum": 500, "default": 50}, + }, + "required": [], + "additionalProperties": False, + }, + }, + { + "name": "getManagedNe", + "description": "Get single managed NE metadata (connect_status, hop config summary).", + "inputSchema": { + "type": "object", + "properties": {"ne_id": {"type": "string"}}, + "required": ["ne_id"], + "additionalProperties": False, + }, + }, + { + "name": "execManagedNe", + "description": "Run read-only CLI on a managed NE via netx (show/display/ping; max 5 commands).", + "inputSchema": { + "type": "object", + "properties": { + "ne_id": {"type": "string"}, + "commands": {"type": "array", "items": {"type": "string"}, "minItems": 1, "maxItems": 5}, + "read_timeout_sec": {"type": "integer", "minimum": 10, "maximum": 120}, + }, + "required": ["ne_id", "commands"], + "additionalProperties": False, + }, + }, +] + +_HANDLERS: dict[str, Callable[[dict[str, Any]], dict[str, Any]]] = { + "queryUmeAlarms": _query_ume_alarms, + "aggregateUmeAlarms": _aggregate_ume_alarms, + "runUmeDiagnostics": _run_ume_diagnostics, + "queryUmeNeInventory": _query_ume_ne_inventory, + "getUmeNe": _get_ume_ne, + "queryUmeAlarmsRaw": _query_ume_alarms_raw, + "aggregateUmeAlarmsRaw": _aggregate_ume_alarms_raw, + "listUmeAlarmFields": _list_ume_alarm_fields, + "sqlQueryUme": _sql_query_ume, + "listManagedNe": _list_managed_ne, + "getManagedNe": _get_managed_ne, + "execManagedNe": _exec_managed_ne, +} + + +def call_http_tool(name: str, args: dict[str, Any]) -> dict[str, Any]: + fn = _HANDLERS.get(str(name or "").strip()) + if not fn: + raise ValueError(f"unknown tool: {name}") + return mcp_from_handler_result(fn(dict(args or {}))) diff --git a/packages/netx-mcp/src/netx_mcp/server.py b/packages/netx-mcp/src/netx_mcp/server.py new file mode 100644 index 0000000..82d13d8 --- /dev/null +++ b/packages/netx-mcp/src/netx_mcp/server.py @@ -0,0 +1,87 @@ +"""netx stdio MCP server (HTTP client to netx REST API). + +Environment: +- ``NETX_API_URL``: netx REST base URL (default ``http://127.0.0.1:8890``) +- ``NETX_API_TOKEN``: optional Bearer token +- ``NETX_LANG``: ``zh`` or ``en`` for localized API responses +""" + +from __future__ import annotations + +import json +import sys +from typing import Any + +from netx_mcp.http_tools import HTTP_MCP_TOOLS, call_http_tool + + +def _ensure_utf8_stdio() -> None: + """MCP stdio must be UTF-8; Windows defaults to GBK and breaks the host reader.""" + for stream in (sys.stdin, sys.stdout, sys.stderr): + if stream is None or not hasattr(stream, "reconfigure"): + continue + try: + stream.reconfigure(encoding="utf-8", errors="replace") + except Exception: + pass + + +def _ok(rid: Any, result: dict[str, Any]) -> None: + sys.stdout.write(json.dumps({"jsonrpc": "2.0", "id": rid, "result": result}, ensure_ascii=False) + "\n") + sys.stdout.flush() + + +def _err(rid: Any, code: int, message: str) -> None: + sys.stdout.write( + json.dumps({"jsonrpc": "2.0", "id": rid, "error": {"code": code, "message": message}}, ensure_ascii=False) + + "\n" + ) + sys.stdout.flush() + + +def run_stdio_loop() -> None: + for line in sys.stdin: + raw = line.strip() + if not raw: + continue + try: + req = json.loads(raw) + except Exception: + continue + rid = req.get("id") + method = str(req.get("method") or "") + params = req.get("params") if isinstance(req.get("params"), dict) else {} + + try: + if method == "initialize": + _ok( + rid, + { + "protocolVersion": "2024-11-05", + "capabilities": {"tools": {}}, + "serverInfo": {"name": "netx-mcp", "version": "0.2.0", "mode": "http"}, + }, + ) + continue + if method == "notifications/initialized": + continue + if method == "tools/list": + _ok(rid, {"tools": HTTP_MCP_TOOLS}) + continue + if method == "tools/call": + name = str(params.get("name") or "") + args = params.get("arguments") if isinstance(params.get("arguments"), dict) else {} + _ok(rid, call_http_tool(name, args)) + continue + _err(rid, -32601, f"method not found: {method}") + except Exception as exc: + _err(rid, -32000, str(exc)) + + +def main() -> None: + _ensure_utf8_stdio() + run_stdio_loop() + + +if __name__ == "__main__": + main() diff --git a/packages/netx-mcp/tests/test_mcp_http.py b/packages/netx-mcp/tests/test_mcp_http.py new file mode 100644 index 0000000..1a2d512 --- /dev/null +++ b/packages/netx-mcp/tests/test_mcp_http.py @@ -0,0 +1,87 @@ +"""Tests for netx HTTP MCP server.""" + +from __future__ import annotations + +import json +import subprocess +import sys +from unittest.mock import patch + +import pytest + +from netx_mcp.http_tools import HTTP_MCP_TOOLS, call_http_tool + + +def test_http_mcp_tool_list_has_twelve_tools() -> None: + names = [str(t.get("name") or "") for t in HTTP_MCP_TOOLS] + assert len(names) == 12 + assert "queryUmeAlarms" in names + assert "queryUmeAlarmsRaw" in names + assert "execManagedNe" in names + + +def test_call_query_ume_alarms_forwards_http() -> None: + with patch("netx_mcp.http_tools.http_json") as mock_http: + mock_http.return_value = {"ok": True, "data": {"total": 0, "items": []}} + out = call_http_tool("queryUmeAlarms", {"severity": "critical", "page": 1, "page_size": 10}) + mock_http.assert_called_once() + assert mock_http.call_args[0][0] == "GET" + assert mock_http.call_args[0][1] == "/v1/ume/alarms" + params = mock_http.call_args[1]["params"] + assert params["severity"] == "critical" + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["ok"] is True + + +def test_call_get_ume_ne_requires_id() -> None: + out = call_http_tool("getUmeNe", {}) + assert out.get("isError") is True + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["error"] == "ne_id_required" + + +def test_call_exec_managed_ne_posts_body() -> None: + with patch("netx_mcp.http_tools.http_post_json") as mock_post: + mock_post.return_value = {"ok": True, "data": {"ok": True, "output": "ok"}} + out = call_http_tool( + "execManagedNe", + {"ne_id": "abc", "commands": ["show version"]}, + ) + mock_post.assert_called_once() + assert mock_post.call_args[0][0] == "/v1/managed-ne/exec" + body = mock_post.call_args[0][1] + assert body["ne_id"] == "abc" + assert body["commands"] == ["show version"] + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["ok"] is True + + +def test_stdio_initialize_and_tools_list() -> None: + proc = subprocess.Popen( + [sys.executable, "-m", "netx_mcp"], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + assert proc.stdin and proc.stdout + init_req = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {}}) + "\n" + proc.stdin.write(init_req) + proc.stdin.flush() + init_line = proc.stdout.readline() + init_resp = json.loads(init_line) + assert init_resp["result"]["serverInfo"]["mode"] == "http" + + list_req = json.dumps({"jsonrpc": "2.0", "id": 2, "method": "tools/list", "params": {}}) + "\n" + proc.stdin.write(list_req) + proc.stdin.flush() + list_line = proc.stdout.readline() + list_resp = json.loads(list_line) + tools = list_resp["result"]["tools"] + assert len(tools) == 12 + + proc.terminate() + proc.wait(timeout=5) diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..9ee6cd5 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,37 @@ +[build-system] +requires = ["setuptools>=68", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "netx-ops" +version = "0.2.0" +description = "netx operations tool: alarm-centric workflows with REST API and stdio MCP" +readme = "README.md" +requires-python = ">=3.11" +license = { text = "MIT" } +dependencies = [ + "fastapi>=0.115.0", + "uvicorn>=0.30.0", + "httpx>=0.27.0", + "sqlalchemy>=2.0.0", + "psycopg[binary]>=3.2.0", + "pandas>=2.0.0", + "openpyxl>=3.1.0", + "PyYAML>=6.0.0", + "pydantic>=2.8.0", + "pydantic-settings>=2.3.0", + "python-multipart>=0.0.9", + "websocket-client>=1.8.0", + "cryptography>=42.0.0", + "netmiko>=4.3.0", +] + +[project.optional-dependencies] +mcp = ["netx-mcp @ file:packages/netx-mcp"] + +[project.scripts] +netx-mcp = "netx_mcp.server:main" + +[tool.setuptools.packages.find] +where = ["."] +include = ["netx_api*"] diff --git a/tests/test_mcp_http.py b/tests/test_mcp_http.py new file mode 100644 index 0000000..4181e97 --- /dev/null +++ b/tests/test_mcp_http.py @@ -0,0 +1,107 @@ +"""Tests for netx HTTP MCP (via netx-mcp package).""" + +from __future__ import annotations + +import json +import subprocess +import sys +from unittest.mock import patch + +import pytest + +from netx_mcp.http_tools import HTTP_MCP_TOOLS, call_http_tool + + +def test_http_mcp_tool_list_has_twelve_tools() -> None: + names = [str(t.get("name") or "") for t in HTTP_MCP_TOOLS] + assert len(names) == 12 + assert "queryUmeAlarms" in names + assert "queryUmeAlarmsRaw" in names + assert "execManagedNe" in names + + +def test_call_query_ume_alarms_forwards_http() -> None: + with patch("netx_mcp.http_tools.http_json") as mock_http: + mock_http.return_value = {"ok": True, "data": {"total": 0, "items": []}} + out = call_http_tool("queryUmeAlarms", {"severity": "critical", "page": 1, "page_size": 10}) + mock_http.assert_called_once() + assert mock_http.call_args[0][0] == "GET" + assert mock_http.call_args[0][1] == "/v1/ume/alarms" + params = mock_http.call_args[1]["params"] + assert params["severity"] == "critical" + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["ok"] is True + + +def test_call_get_ume_ne_requires_id() -> None: + out = call_http_tool("getUmeNe", {}) + assert out.get("isError") is True + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["error"] == "ne_id_required" + + +def test_call_exec_managed_ne_posts_body() -> None: + with patch("netx_mcp.http_tools.http_post_json") as mock_post: + mock_post.return_value = {"ok": True, "data": {"ok": True, "output": "ok"}} + out = call_http_tool( + "execManagedNe", + {"ne_id": "abc", "commands": ["show version"]}, + ) + mock_post.assert_called_once() + assert mock_post.call_args[0][0] == "/v1/managed-ne/exec" + body = mock_post.call_args[0][1] + assert body["ne_id"] == "abc" + assert body["commands"] == ["show version"] + text = out["content"][0]["text"] + payload = json.loads(text) + assert payload["ok"] is True + + +def test_stdio_initialize_and_tools_list() -> None: + proc = subprocess.Popen( + [sys.executable, "-m", "netx_mcp"], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + assert proc.stdin and proc.stdout + init_req = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {}}) + "\n" + proc.stdin.write(init_req) + proc.stdin.flush() + init_line = proc.stdout.readline() + init_resp = json.loads(init_line) + assert init_resp["result"]["serverInfo"]["mode"] == "http" + + list_req = json.dumps({"jsonrpc": "2.0", "id": 2, "method": "tools/list", "params": {}}) + "\n" + proc.stdin.write(list_req) + proc.stdin.flush() + list_line = proc.stdout.readline() + list_resp = json.loads(list_line) + tools = list_resp["result"]["tools"] + assert len(tools) == 12 + + proc.terminate() + proc.wait(timeout=5) + + +def test_legacy_netx_api_mcp_module_still_works() -> None: + proc = subprocess.Popen( + [sys.executable, "-m", "netx_api.mcp"], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + cwd=str(__import__("pathlib").Path(__file__).resolve().parents[1]), + ) + assert proc.stdin and proc.stdout + init_req = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {}}) + "\n" + proc.stdin.write(init_req) + proc.stdin.flush() + init_line = proc.stdout.readline() + init_resp = json.loads(init_line) + assert init_resp["result"]["serverInfo"]["mode"] == "http" + proc.terminate() + proc.wait(timeout=5)