Ship Cursor-only MCP admin and modular Admin/Chat UI.

Plugins/Skills install and edit via mcpServers JSON, with clearer row actions, soft reloads, and instant loading placeholders; also drop MCP market and harden related scheduler/logging/exec guards.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-12 19:38:26 +08:00
parent 2aff012b48
commit f6f129931f
61 changed files with 18057 additions and 16327 deletions

View file

@ -9,40 +9,41 @@ _BOT_MENTION_RE = re.compile(r"@\S+")
# intent -> (en hint, zh hint)
_HINTS: dict[str, tuple[str, str]] = {
"excel_export": (
"[Ops short-intent: Excel export. Call ume_alarm_xlsx_report(..., deliverable=true) now. "
"write_xlsx / inventory / CLI / run_command are hidden this turn — do not build xlsx via shell.]",
"[短指令:导出 Excel。立即 ume_alarm_xlsx_report(..., deliverable=true)。"
"本轮已隐藏 write_xlsx/清单/CLI/run_command,禁止用 shell 搓表。]",
),
"license": (
"[Ops short-intent: license/capacity. Call ume_alarm_xlsx_report(mode=list, keyword=license, deliverable=true) "
"or aggregateUmeAlarms/queryUmeAlarmsRaw; write_xlsx/CLI/inventory are hidden this turn.]",
"[短指令:License/容量。立即 ume_alarm_xlsx_report(mode=list, keyword=license, deliverable=true) "
"或 aggregate/queryUmeAlarmsRaw;本轮已隐藏 write_xlsx/CLI/清单。]",
),
"congestion": (
"[Ops short-intent: bandwidth congestion. Call ume_alarm_xlsx_report(mode=list, deliverable=true) "
"or aggregateUmeAlarms/queryUmeAlarmsRaw; write_xlsx/CLI/inventory/sql are hidden this turn.]",
"[短指令:带宽拥塞。立即 ume_alarm_xlsx_report(mode=list, deliverable=true) 或 aggregate/query;"
"本轮已隐藏 write_xlsx/CLI/清单/sql。]",
),
"fiber_cut": (
"[Ops short-intent: fiber/LOS. Call ume_alarm_xlsx_report(mode=fiber_cut, deliverable=true) now. "
"Inventory/CLI tools are hidden this turn — do not try listCliTargets/execManagedNe.]",
"write_xlsx / inventory / CLI are hidden this turn — do not try listCliTargets/execManagedNe/write_xlsx.]",
"[短指令:断纤/LOS。立即 ume_alarm_xlsx_report(mode=fiber_cut, deliverable=true)。"
"本轮已隐藏清单/CLI 工具,勿调用 listCliTargets/execManagedNe。]",
"本轮已隐藏 write_xlsx/清单/CLI,勿调用 listCliTargets/execManagedNe/write_xlsx。]",
),
"offline": (
"[Ops short-intent: offline NE. Call ume_alarm_xlsx_report(mode=offline, deliverable=true) now. "
"Inventory/CLI tools are hidden this turn.]",
"[短指令:离线网元。立即 ume_alarm_xlsx_report(mode=offline, deliverable=true)。本轮已隐藏清单/CLI。]",
"write_xlsx / inventory / CLI are hidden this turn.]",
"[短指令:离线网元。立即 ume_alarm_xlsx_report(mode=offline, deliverable=true)。"
"本轮已隐藏 write_xlsx/清单/CLI。]",
),
"alarm_tally": (
"[Ops short-intent: alarm tally/top. Prefer ume_alarm_xlsx_report(mode=aggregate_by_host) or "
"aggregateUmeAlarms; inventory/CLI tools are hidden this turn.]",
"[短指令:告警统计/Top。优先 ume_alarm_xlsx_report(mode=aggregate_by_host) 或 aggregateUmeAlarms;"
"本轮已隐藏清单/CLI。]",
),
"excel_export": (
"[Ops short-intent: Excel export. Prefer ume_alarm_xlsx_report or write_xlsx(deliverable=true). "
"Inventory/CLI/run_command are hidden this turn — do not build xlsx via shell.]",
"[短指令:导出 Excel。优先 ume_alarm_xlsx_report 或 write_xlsx(deliverable=true);"
"本轮已隐藏清单/CLI/run_command。]",
),
"license": (
"[Ops short-intent: license/capacity. Prefer ume_alarm_xlsx_report(mode=list, keyword=license) or "
"aggregateUmeAlarms/queryUmeAlarmsRaw; CLI/inventory tools are hidden this turn.]",
"[短指令:License/容量。优先 ume_alarm_xlsx_report(mode=list, keyword=license) 或 "
"aggregate/queryUmeAlarmsRaw;本轮已隐藏 CLI/清单。]",
),
"congestion": (
"[Ops short-intent: bandwidth congestion. Prefer ume_alarm_xlsx_report(mode=list) or "
"aggregateUmeAlarms/queryUmeAlarmsRaw; CLI/inventory/sql are hidden this turn.]",
"[短指令:带宽拥塞。优先 ume_alarm_xlsx_report(mode=list) 或 aggregate/query;"
"本轮已隐藏 CLI/清单/sql。]",
"[Ops short-intent: alarm tally/top. Call ume_alarm_xlsx_report(mode=aggregate_by_host, deliverable=true) "
"or aggregateUmeAlarms; write_xlsx / inventory / CLI are hidden this turn.]",
"[短指令:告警统计/Top。立即 ume_alarm_xlsx_report(mode=aggregate_by_host, deliverable=true) "
"或 aggregateUmeAlarms;本轮已隐藏 write_xlsx/清单/CLI。]",
),
"continue": (
"[Ops short-intent: continue/confirm. Resume the unfinished prior task immediately; "
@ -68,6 +69,7 @@ _SUPPRESSED_TOOL_NAMES = frozenset(
"findtopologypaths",
"sqlqueryume",
"run_command",
"write_xlsx",
"netx_list_managed_ne",
"netx_get_managed_ne",
"netx_exec_managed_ne",

View file

@ -0,0 +1,182 @@
"""Guards for mcp__netx__execManagedNe / netx_exec_managed_ne spam on field turns."""
from __future__ import annotations
import json
import os
from typing import Any
def _env_int(name: str, default: int, *, min_v: int = 1, max_v: int = 50) -> int:
raw = str(os.getenv(name) or "").strip()
if not raw:
return default
try:
n = int(raw)
except Exception:
return default
return max(min_v, min(int(n), max_v))
def exec_managed_ne_single_budget() -> int:
"""Max single-NE execManagedNe calls per turn before requiring ne_ids/ume_ne_ids batch."""
return _env_int("AIA_EXEC_MANAGED_NE_SINGLE_BUDGET", 6, min_v=2, max_v=30)
def exec_managed_ne_fail_budget() -> int:
"""Max failed execManagedNe (any mode) per turn before blocking further single-NE calls."""
return _env_int("AIA_EXEC_MANAGED_NE_FAIL_BUDGET", 5, min_v=2, max_v=30)
def is_exec_managed_ne_tool(name: str) -> bool:
raw = str(name or "").strip()
if not raw:
return False
if "__" in raw:
raw = raw.rsplit("__", 1)[-1]
key = raw.strip().lower().replace("-", "_")
return key in {"execmanagedne", "netx_exec_managed_ne", "exec_managed_ne"}
def is_batch_exec_args(args: dict[str, Any] | None) -> bool:
a = args if isinstance(args, dict) else {}
for key in ("ne_ids", "ume_ne_ids", "targets"):
val = a.get(key)
if isinstance(val, list) and len(val) > 0:
return True
return False
def normalize_exec_managed_ne_args(args: dict[str, Any] | None) -> dict[str, Any]:
"""Default / clamp read_timeout_sec so agents stop hitting 30s walls."""
out = dict(args or {})
rts = out.get("read_timeout_sec")
if rts is None or str(rts).strip() == "":
out["read_timeout_sec"] = 60
else:
try:
out["read_timeout_sec"] = max(10, min(120, int(rts)))
except Exception:
out["read_timeout_sec"] = 60
return out
def _parse_json_obj(raw: Any) -> dict[str, Any]:
if isinstance(raw, dict):
return dict(raw)
if isinstance(raw, str) and raw.strip():
try:
data = json.loads(raw)
return data if isinstance(data, dict) else {}
except Exception:
return {}
return {}
def load_turn_exec_managed_ne_stats(
store: Any,
*,
session_id: str,
turn_uuid: str,
) -> tuple[int, int, int]:
"""Return (single_calls, batch_calls, fail_calls) already persisted this turn."""
tu = str(turn_uuid or "").strip()
sid = str(session_id or "").strip()
single = batch = fails = 0
if not tu or not sid:
return single, batch, fails
try:
rows = store.get_messages(session_id=sid, limit=500)
except Exception:
return single, batch, fails
for m in rows or []:
if str(getattr(m, "role", "") or "").strip().lower() != "tool":
continue
if str(getattr(m, "turn_uuid", "") or "").strip() != tu:
continue
ep = _parse_json_obj(getattr(m, "event_payload", None))
name = str(ep.get("tool_name") or "").strip()
if not name:
raw_tc = getattr(m, "tool_calls", None)
tc = _parse_json_obj(raw_tc)
name = str(tc.get("name") or "").strip()
if not is_exec_managed_ne_tool(name):
continue
mode = str(ep.get("exec_ne_mode") or "").strip().lower()
if mode == "batch":
batch += 1
else:
# Missing mode (older rows) → treat as single (conservative).
single += 1
if ep.get("ok") is False:
fails += 1
continue
try:
payload = json.loads(str(getattr(m, "content", "") or "") or "{}")
except Exception:
payload = {}
if isinstance(payload, dict) and payload.get("ok") is False:
fails += 1
return single, batch, fails
def budget_block_payload(
*,
reason: str,
lang: str = "en",
single_used: int,
single_budget: int,
fail_used: int = 0,
fail_budget: int = 0,
) -> dict[str, Any]:
en = str(lang or "").strip().lower().startswith("en")
if reason == "fail_budget":
code = "cli_fail_budget_exceeded"
hint = (
f"execManagedNe already failed {fail_used}/{fail_budget} times this turn. "
"Stop one-NE loops; use one execManagedNe(ne_ids|ume_ne_ids=..., commands=...) batch "
"or summarize reachable failures — do not keep probing."
if en
else f"本轮 execManagedNe 已失败 {fail_used}/{fail_budget} 次。"
"停止单台循环;改用一次 ne_ids/ume_ne_ids 批量,或汇总可达性失败,勿继续盲探。"
)
err = "cli_fail_budget_exceeded"
else:
code = "cli_call_budget_exceeded"
hint = (
f"Single-NE execManagedNe budget exhausted ({single_used}/{single_budget} this turn). "
"For more NEs call ONE execManagedNe with ne_ids[] or ume_ne_ids[] (shared commands). "
"Do not loop one-NE execManagedNe."
if en
else f"单台 execManagedNe 预算已用尽(本轮 {single_used}/{single_budget})。"
"更多网元请一次传入 ne_ids[] / ume_ne_ids[] 批量执行,禁止逐台循环。"
)
err = "cli_call_budget_exceeded"
return {
"ok": False,
"error_code": code,
"failure_class": "retry_guard",
"error": err,
"hint": hint,
"example": {
"ume_ne_ids": ["<id1>", "<id2>", "<id3>"],
"commands": ["show version"],
"read_timeout_sec": 90,
"concurrency": 4,
},
"single_used": int(single_used),
"single_budget": int(single_budget),
"fail_used": int(fail_used),
"fail_budget": int(fail_budget),
}
__all__ = [
"budget_block_payload",
"exec_managed_ne_fail_budget",
"exec_managed_ne_single_budget",
"is_batch_exec_args",
"is_exec_managed_ne_tool",
"load_turn_exec_managed_ne_stats",
"normalize_exec_managed_ne_args",
]

View file

@ -1087,6 +1087,24 @@ class ToolExecutor:
session_id=ctx.session_id,
turn_uuid=str(ctx.turn_uuid or ""),
)
from runtime.chat.exec_managed_ne_guard import (
budget_block_payload,
exec_managed_ne_fail_budget,
exec_managed_ne_single_budget,
is_batch_exec_args,
is_exec_managed_ne_tool,
load_turn_exec_managed_ne_stats,
)
prior_single_exec, _prior_batch_exec, prior_exec_fails = load_turn_exec_managed_ne_stats(
ctx.store,
session_id=ctx.session_id,
turn_uuid=str(ctx.turn_uuid or ""),
)
single_exec_budget = exec_managed_ne_single_budget()
fail_exec_budget = exec_managed_ne_fail_budget()
local_single_exec = 0
local_exec_fails = 0
results_by_id: dict[str, tuple[dict[str, Any], int]] = {}
runnable_tool_uses: list[LLMToolCall] = []
@ -1216,6 +1234,48 @@ class ToolExecutor:
},
)
continue
if is_exec_managed_ne_tool(tool_name):
batchish = is_batch_exec_args(dict(tc.arguments or {}))
single_used = int(prior_single_exec) + int(local_single_exec)
fail_used = int(prior_exec_fails) + int(local_exec_fails)
if (not batchish) and fail_used >= int(fail_exec_budget):
results_by_id[tc.id] = (
budget_block_payload(
reason="fail_budget",
lang=str(ctx.lang or "en"),
single_used=single_used,
single_budget=single_exec_budget,
fail_used=fail_used,
fail_budget=fail_exec_budget,
),
0,
)
_trace(
"cli_fail_budget_exceeded",
{"tool_name": tc.name, "fail_used": fail_used, "fail_budget": fail_exec_budget},
)
continue
if (not batchish) and single_used >= int(single_exec_budget):
results_by_id[tc.id] = (
budget_block_payload(
reason="call_budget",
lang=str(ctx.lang or "en"),
single_used=single_used,
single_budget=single_exec_budget,
fail_used=fail_used,
fail_budget=fail_exec_budget,
),
0,
)
_trace(
"cli_call_budget_exceeded",
{
"tool_name": tc.name,
"single_used": single_used,
"single_budget": single_exec_budget,
},
)
continue
count = int(sig_seen.get(sig, 0))
name_low = str(tc.name or "").strip().lower()
listish = name_low.endswith(
@ -1270,6 +1330,8 @@ class ToolExecutor:
)
continue
first_tool_call_id_by_signature[sig] = str(tc.id or "")
if is_exec_managed_ne_tool(tool_name) and not is_batch_exec_args(dict(tc.arguments or {})):
local_single_exec += 1
runnable_tool_uses.append(tc)
for batch in partition_tool_use_batches(runnable_tool_uses, ctx.tools):
@ -1332,6 +1394,16 @@ class ToolExecutor:
result = normalize_tool_result(result)
if isinstance(result, dict) and result.get("ok") is False:
failed_signatures.add(f"{tc.name}:{self._json_dumps_safe(dict(tc.arguments or {}))}")
if is_exec_managed_ne_tool(str(tc.name or "")):
# Count blocked budget responses too so fail-budget can engage same turn.
if str(result.get("error_code") or "") not in {
"cli_call_budget_exceeded",
"cli_fail_budget_exceeded",
"identical_retry_blocked",
"retry_forbidden_blocked",
"tool_loop_guard",
}:
local_exec_fails += 1
if isinstance(result, dict) and _result_is_retry_forbidden(result):
retry_forbidden_tools.add(str(tc.name or ""))
persisted_result, ingested_refs = ingest_embedded_image_blobs_as_refs(
@ -1376,6 +1448,9 @@ class ToolExecutor:
tool_content = self._json_dumps_safe(result_for_llm)
t_db2 = time.perf_counter()
tool_sig = f"{tc.name}:{self._json_dumps_safe(dict(tc.arguments or {}))}"
exec_ne_mode = ""
if is_exec_managed_ne_tool(str(tc.name or "")):
exec_ne_mode = "batch" if is_batch_exec_args(dict(tc.arguments or {})) else "single"
msg_row = ctx.store.add_message(
session_id=ctx.session_id,
role="tool",
@ -1393,6 +1468,7 @@ class ToolExecutor:
"retry_forbidden": bool(
isinstance(result, dict) and _result_is_retry_forbidden(result)
),
**({"exec_ne_mode": exec_ne_mode} if exec_ne_mode else {}),
},
)
try:

View file

@ -98,6 +98,8 @@ class TurnIdleTracker:
"identical_retry_blocked",
"retry_forbidden_blocked",
"tool_loop_guard",
"cli_call_budget_exceeded",
"cli_fail_budget_exceeded",
}:
guard += 1
stats = RoundStats(

View file

@ -1,4 +1,4 @@
"""Build a reinstallable MCP JSON from SQLite and persist under ``oclaw/_local/`` for migration."""
"""Build a reinstallable Cursor ``mcpServers`` JSON from SQLite under ``oclaw/_local/``."""
from __future__ import annotations
import json
@ -10,6 +10,7 @@ from typing import Any
from svc.config.paths import PROJECT_ROOT
from svc.persistence.sqlite_store import SqliteStore
from runtime.tools.mcp.cursor_config import build_cursor_mcp_export
from runtime.tools.mcp.registry import McpRegistry
_EXPORT_FILENAME = "mcp_registry_migrated.json"
@ -22,11 +23,7 @@ def mcp_migrated_json_path() -> Path:
def _collapse_entry_arg(arg: str, root: Path) -> str:
"""Rewrite absolute / repo-relative paths to ``__REPO_ROOT__/...`` when they exist on disk under root.
Skips npx flags (``-y``) and tokens that are not an existing file or directory, so package names
(e.g. ``mcp-sqlite``, ``@upstash/foo``) are not mistaken for paths.
"""
"""Rewrite absolute / repo-relative paths to ``__REPO_ROOT__/...`` when under root."""
t = (arg or "").strip()
if not t or t.startswith("__REPO_ROOT__") or t.startswith("-"):
return str(arg)
@ -52,41 +49,24 @@ def _collapse_entry_args(args: list[str], root: Path) -> list[str]:
def build_mcp_install_export_document(store: SqliteStore) -> dict[str, Any]:
rows = McpRegistry(store).list_servers(enabled_only=False)
root = PROJECT_ROOT.resolve()
servers: list[dict[str, Any]] = []
collapsed: list[dict[str, Any]] = []
for r in rows:
if not isinstance(r, dict):
continue
raw_args = r.get("entry_args")
row = dict(r)
raw_args = row.get("entry_args")
if isinstance(raw_args, list):
entry_args = _collapse_entry_args([str(x) for x in raw_args if str(x).strip()], root)
row["entry_args"] = _collapse_entry_args([str(x) for x in raw_args if str(x).strip()], root)
else:
entry_args = []
es = r.get("env_schema")
env_schema = es if isinstance(es, dict) else {}
perms = r.get("required_permissions")
if not isinstance(perms, list):
perms = []
servers.append(
{
"server_id": str(r.get("server_id") or ""),
"source_type": str(r.get("source_type") or ""),
"source_ref": str(r.get("source_ref") or ""),
"version": str(r.get("version") or ""),
"entry_command": str(r.get("entry_command") or ""),
"entry_args": entry_args,
"env_schema": env_schema,
"required_permissions": [str(x) for x in perms if str(x).strip()],
"risk_level": str(r.get("risk_level") or "high"),
"enabled": bool(r.get("enabled")),
"timeout_s": float(r.get("timeout_s") or 30.0),
"dry_run": False,
}
)
return {
"_comment": "管理台导出 / 新安装后自动落盘。可粘到「Install from JSON」;`__REPO_ROOT__/` 在管理台与 seed 中展开为仓库根。已存在路径会写成占位符,其余参数保持原样。",
"exported_at": datetime.now(timezone.utc).isoformat(),
"servers": servers,
}
row["entry_args"] = []
collapsed.append(row)
doc = build_cursor_mcp_export(collapsed)
doc["_comment"] = (
"Cursor mcpServers export. Paste into Admin → Plugins → Install, "
"or feed to seed_mcp_registry.py. Paths under the repo may use __REPO_ROOT__/."
)
doc["exported_at"] = datetime.now(timezone.utc).isoformat()
return doc
def persist_mcp_migrated_file(store: SqliteStore) -> str | None:

View file

@ -1,10 +1,10 @@
"""将 oclaw/data/mcp_registry.seed.json 中的 MCP 定义写入本地 SQLite(与管理台 Install 等价)。
"""将 Cursor ``mcpServers`` JSON 写入本地 SQLite(与管理台 Install 等价)。
库表被清空、换库或新环境时:在项目根执行
python oclaw/scripts/seed_mcp_registry.py
可选:python oclaw/scripts/seed_mcp_registry.py path/to/other.seed.json
python -m runtime.operations.scripts.seed_mcp_registry
可选:python -m runtime.operations.scripts.seed_mcp_registry path/to/mcp.json
随后在各 MCP 上点 Health → Sync Tools;Context7 等需在 src/_local/mcp_local.env 配密钥。
随后在各 MCP 上点 Health → Sync Tools;密钥放在 oclaw/_local/mcp_local.env。
"""
from __future__ import annotations
@ -12,13 +12,15 @@ import json
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
ROOT = Path(__file__).resolve().parents[3]
sys.path.insert(0, str(ROOT))
from svc.config.paths import PROJECT_ROOT, db_path # noqa: E402
from svc.persistence.sqlite_store import SqliteStore # noqa: E402
from svc.persistence.assistant_store import get_assistant_store
from runtime.tools.mcp.installer import McpServerManifest, _safe_server_id, install_mcp_server # noqa: E402
from svc.config.paths import PROJECT_ROOT # noqa: E402
from svc.persistence.assistant_store import get_assistant_store # noqa: E402
from runtime.tools.mcp.cursor_config import parse_cursor_mcp_document # noqa: E402
from runtime.tools.mcp.installer import install_mcp_server # noqa: E402
from runtime.tools.mcp.manifest import McpServerManifest # noqa: E402
from runtime.operations.mcp_registry_export import persist_mcp_migrated_file # noqa: E402
def _subst_repo(path_str: str) -> str:
@ -34,37 +36,30 @@ def main(argv: list[str]) -> int:
print(f"seed file not found: {seed_path}", file=sys.stderr)
return 2
raw = json.loads(seed_path.read_text(encoding="utf-8"))
items = raw.get("servers") if isinstance(raw, dict) else raw
if not isinstance(items, list) or not items:
print("no servers[] in seed file", file=sys.stderr)
try:
items = parse_cursor_mcp_document(raw)
except ValueError as exc:
print(f"invalid Cursor mcpServers document: {exc}", file=sys.stderr)
return 2
store = get_assistant_store()
ok_n = 0
for payload in items:
if not isinstance(payload, dict):
continue
source_type = str(payload.get("source_type") or "").strip().lower()
source_ref = str(payload.get("source_ref") or "").strip()
if source_type not in {"github", "npm", "pypi"} or not source_ref:
print(f"skip invalid: {payload.get('server_id')!r}")
continue
server_id = _safe_server_id(str(payload.get("server_id") or source_ref))
entry_args = [_subst_repo(x) for x in (payload.get("entry_args") or []) if str(x).strip()]
dry_run = bool(payload.get("dry_run", False))
manifest = McpServerManifest(
server_id=server_id,
source_type=source_type,
source_ref=source_ref,
server_id=str(payload.get("server_id") or ""),
source_type=str(payload.get("source_type") or "local"),
source_ref=str(payload.get("source_ref") or ""),
version=str(payload.get("version") or "").strip(),
entry_command=str(payload.get("entry_command") or "").strip(),
entry_args=entry_args,
env_schema=payload.get("env_schema") if isinstance(payload.get("env_schema"), dict) else {},
permissions=[str(x) for x in (payload.get("required_permissions") or [])],
risk_level=str(payload.get("risk_level") or "high"),
enabled=bool(payload.get("enabled")),
enabled=bool(payload.get("enabled", True)),
timeout_s=float(payload.get("timeout_s") or 30.0),
)
dry_run = bool(payload.get("dry_run", False))
inst = install_mcp_server(manifest, dry_run=dry_run)
store.upsert_mcp_server(
server_id=manifest.server_id,
@ -86,10 +81,14 @@ def main(argv: list[str]) -> int:
detail={"error": inst.error, **(inst.details or {}), "seed": True},
install_command=inst.install_command,
)
print(f"{server_id}: install_ok={inst.ok} enabled={manifest.enabled if inst.ok else False} cmd={inst.install_command[:120]!r}")
print(
f"{manifest.server_id}: install_ok={inst.ok} "
f"enabled={manifest.enabled if inst.ok else False} cmd={inst.install_command[:120]!r}"
)
if inst.ok:
ok_n += 1
print(f"done: {ok_n}/{len(items)} npm/pypi install steps succeeded; rows upserted for all valid entries.")
persist_mcp_migrated_file(store)
print(f"done: {ok_n}/{len(items)} install steps succeeded; rows upserted for all valid entries.")
return 0

View file

@ -0,0 +1,65 @@
"""Classify scheduled_job_run errors for Admin dashboards and ops triage."""
from __future__ import annotations
from typing import Any
def classify_scheduled_job_error(error: str | None, *, status: str | None = None) -> str:
"""Return a coarse failure_class for a scheduled job run.
Classes: overlap | timeout | delivery | mcp | auth | cancelled | runtime | "" (success/empty).
"""
st = str(status or "").strip().lower()
err = str(error or "").strip()
if st in {"success", "ok", "skipped"} and not err:
return ""
if st == "skipped" or "overlapping" in err.lower() or err.lower() == "overlapping_run":
return "overlap"
if not err and st not in {"failed", "error"}:
return ""
blob = err.lower()
if "stale_running" in blob:
return "stale"
if any(x in blob for x in ("timeout", "timed out", "deadline", "read_timeout")):
return "timeout"
if any(
x in blob
for x in (
"delivery",
"enqueue_failed",
"connection closed",
"whatsapp",
"weixin",
"send_failed",
"outbound",
)
):
return "delivery"
if any(x in blob for x in ("insufficient_scope", "unauthorized", "forbidden", "permission denied", "auth")):
return "auth"
if any(x in blob for x in ("mcp", "unregistered tool", "tool_not_registered")):
return "mcp"
if any(x in blob for x in ("cancel", "interrupted", "stopped")):
return "cancelled"
if err:
return "runtime"
if st in {"failed", "error"}:
return "runtime"
return ""
def enrich_scheduled_job_run_dict(d: dict[str, Any]) -> dict[str, Any]:
out = dict(d or {})
fc = classify_scheduled_job_error(out.get("error"), status=str(out.get("status") or ""))
if fc:
out["failure_class"] = fc
elif str(out.get("status") or "").strip().lower() in {"failed", "error"}:
out["failure_class"] = "runtime"
return out
__all__ = [
"classify_scheduled_job_error",
"enrich_scheduled_job_run_dict",
]

View file

@ -1,16 +1,25 @@
from .manifest import McpServerManifest
from .installer import McpInstallResult, install_mcp_server
from .runtime import McpProcessRuntime
from __future__ import annotations
from .adapter import materialize_mcp_tools
from .cursor_config import (
build_cursor_mcp_export,
parse_cursor_mcp_document,
registry_row_to_cursor_server,
)
from .installer import McpInstallResult, install_mcp_server, uninstall_mcp_server
from .manifest import McpServerManifest
from .registry import McpRegistry
from .market import search_mcp_market
from .runtime import McpProcessRuntime
__all__ = [
"McpServerManifest",
"McpInstallResult",
"McpProcessRuntime",
"McpRegistry",
"McpServerManifest",
"build_cursor_mcp_export",
"install_mcp_server",
"materialize_mcp_tools",
"McpRegistry",
"search_mcp_market",
"parse_cursor_mcp_document",
"registry_row_to_cursor_server",
"uninstall_mcp_server",
]

View file

@ -150,6 +150,13 @@ class _McpBoundTool:
def _handler(args: dict[str, Any]) -> dict[str, Any]:
call_args = dict(args or {})
from runtime.chat.exec_managed_ne_guard import (
is_exec_managed_ne_tool,
normalize_exec_managed_ne_args,
)
if is_exec_managed_ne_tool(tool_name):
call_args = normalize_exec_managed_ne_args(call_args)
cache_ttl = _mcp_list_cache_ttl(tool_name)
cache_key = ""
if cache_ttl is not None:
@ -180,7 +187,7 @@ class _McpBoundTool:
res["hint"] = (
f"Cache {tool_name} results briefly; reuse ids/rows instead of listing again in the same turn."
)
if tool_name == "execManagedNe" and res.get("ok") is False:
if is_exec_managed_ne_tool(tool_name) and res.get("ok") is False:
from runtime.tools.tool_error_hints import enrich_exec_managed_ne_error
res = enrich_exec_managed_ne_error(res)

View file

@ -0,0 +1,313 @@
# -*- coding: utf-8 -*-
"""Cursor-style ``mcpServers`` JSON ↔ oclaw MCP registry payloads."""
from __future__ import annotations
import json
import re
from typing import Any
_ENV_VAR_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}")
_REMOTE_TYPES = {
"streamablehttp",
"streamable_http",
"streamable-http",
"http",
"sse",
}
def _safe_server_id(seed: str) -> str:
v = re.sub(r"[^a-zA-Z0-9._-]+", "-", str(seed or "").strip().lower()).strip("-")
return v or "mcp-server"
def _env_schema_from_cursor_env(env: dict[str, Any] | None) -> dict[str, Any]:
out: dict[str, Any] = {}
if not isinstance(env, dict):
return out
for ek, ev in env.items():
name = str(ek or "").strip()
if not name:
continue
out[name] = {
"type": "string",
"default": "" if ev is None else str(ev),
"description": "From mcpServers env",
}
return out
def _env_schema_from_header_placeholders(headers: dict[str, Any]) -> dict[str, Any]:
out: dict[str, Any] = {}
for hk, hv in headers.items():
for m in _ENV_VAR_RE.finditer(str(hv or "")):
env_name = str(m.group(1) or "").strip()
if not env_name or env_name in out:
continue
out[env_name] = {
"required": True,
"type": "string",
"description": f"Auto-detected from header {str(hk or '').strip()}",
}
return out
def _cursor_env_from_schema(env_schema: dict[str, Any] | None) -> dict[str, str]:
out: dict[str, str] = {}
if not isinstance(env_schema, dict):
return out
for ek, spec in env_schema.items():
name = str(ek or "").strip()
if not name:
continue
if isinstance(spec, dict):
default = spec.get("default")
out[name] = "" if default is None else str(default)
else:
out[name] = str(spec)
return out
def _is_remote_server(server: dict[str, Any]) -> bool:
transport = str(server.get("type") or "").strip().lower().replace(" ", "")
url = str(server.get("url") or server.get("baseUrl") or "").strip()
if transport in _REMOTE_TYPES:
return True
# Cursor often omits type and only sets url for remote servers.
if url and not str(server.get("command") or server.get("entry_command") or "").strip():
return True
return False
def cursor_server_to_install_payload(server_id: str, server: dict[str, Any] | None) -> dict[str, Any]:
"""Convert one Cursor mcpServers entry into an oclaw install/upsert payload."""
sid = _safe_server_id(server_id)
s = server if isinstance(server, dict) else {}
headers = s.get("headers") if isinstance(s.get("headers"), dict) else {}
env_schema = _env_schema_from_header_placeholders(headers)
enabled = bool(s["isActive"]) if "isActive" in s else True
if "enabled" in s:
enabled = bool(s.get("enabled"))
timeout_s = float(s.get("timeout_s") or 30.0)
if _is_remote_server(s):
base_url = str(s.get("url") or s.get("baseUrl") or "").strip()
if not base_url:
raise ValueError(f"remote_mcp_missing_url:{sid}")
entry_args = ["-y", "mcp-remote", base_url]
for hk, hv in headers.items():
hn = str(hk or "").strip()
hvs = str(hv or "").strip()
if not hn or not hvs:
continue
entry_args.extend(["--header", f"{hn}: {hvs}"])
return {
"server_id": sid,
"source_type": "npm",
"source_ref": "mcp-remote",
"version": "",
"entry_command": "npx",
"entry_args": entry_args,
"env_schema": env_schema,
"required_permissions": [],
"risk_level": "high",
"enabled": enabled,
"timeout_s": timeout_s,
}
cmd = str(s.get("command") or s.get("entry_command") or "").strip()
if not cmd:
raise ValueError(f"stdio_mcp_missing_command:{sid}")
args_raw = s.get("args") if isinstance(s.get("args"), list) else s.get("entry_args")
args = [str(x) for x in (args_raw or [])]
env_from_cursor = _env_schema_from_cursor_env(s.get("env") if isinstance(s.get("env"), dict) else {})
merged = {**env_from_cursor, **env_schema}
if isinstance(s.get("env_schema"), dict):
merged.update(s["env_schema"])
return {
"server_id": sid,
"source_type": "local",
"source_ref": sid,
"version": "",
"entry_command": cmd,
"entry_args": args,
"env_schema": merged,
"required_permissions": [],
"risk_level": "high",
"enabled": enabled,
"timeout_s": timeout_s,
}
def parse_cursor_mcp_document(doc: Any) -> list[dict[str, Any]]:
"""Parse Cursor ``{ mcpServers: {...} }`` (or bare mcpServers object) into install payloads."""
if not isinstance(doc, dict):
raise ValueError("mcp_document_must_be_object")
servers_obj = doc.get("mcpServers")
if servers_obj is None and all(isinstance(v, dict) for v in doc.values()) and doc:
# Allow pasting the inner map directly if every value looks like a server config.
if any(k in next(iter(doc.values()), {}) for k in ("command", "url", "args", "type")):
servers_obj = doc
if not isinstance(servers_obj, dict) or not servers_obj:
raise ValueError("mcpServers_required")
out: list[dict[str, Any]] = []
for raw_key, raw_server in servers_obj.items():
key = str(raw_key or "").strip()
if not key:
continue
out.append(cursor_server_to_install_payload(key, raw_server if isinstance(raw_server, dict) else {}))
if not out:
raise ValueError("mcpServers_empty")
return out
def _try_parse_mcp_remote(entry_command: str, entry_args: list[str]) -> dict[str, Any] | None:
cmd = str(entry_command or "").strip().lower()
args = [str(x) for x in (entry_args or [])]
if cmd not in {"npx", "npx.cmd"} and not cmd.endswith("npx") and not cmd.endswith("npx.cmd"):
# Still allow if args contain mcp-remote (custom launcher).
if "mcp-remote" not in args:
return None
if "mcp-remote" not in args:
return None
url = ""
headers: dict[str, str] = {}
i = 0
while i < len(args):
a = args[i]
if a == "mcp-remote":
i += 1
continue
if a in {"-y", "--yes"}:
i += 1
continue
if a == "--header" and i + 1 < len(args):
hv = args[i + 1]
if ":" in hv:
hk, hval = hv.split(":", 1)
headers[hk.strip()] = hval.strip()
i += 2
continue
if not url and not a.startswith("-"):
url = a
i += 1
continue
i += 1
if not url:
return None
out: dict[str, Any] = {"url": url}
if headers:
out["headers"] = headers
return out
def registry_row_to_cursor_server(row: dict[str, Any]) -> dict[str, Any]:
"""Convert one registry row into a Cursor mcpServers value."""
entry_command = str(row.get("entry_command") or "").strip()
raw_args = row.get("entry_args")
entry_args = [str(x) for x in raw_args] if isinstance(raw_args, list) else []
remote = _try_parse_mcp_remote(entry_command, entry_args)
env = _cursor_env_from_schema(row.get("env_schema") if isinstance(row.get("env_schema"), dict) else {})
if remote is not None:
if env:
# Keep env for header placeholders / runtime.
remote = dict(remote)
remote["env"] = env
return remote
out: dict[str, Any] = {"command": entry_command or "npx", "args": entry_args}
if env:
out["env"] = env
return out
def build_cursor_mcp_export(rows: list[dict[str, Any]]) -> dict[str, Any]:
servers: dict[str, Any] = {}
for row in rows:
if not isinstance(row, dict):
continue
sid = str(row.get("server_id") or "").strip()
if not sid:
continue
servers[sid] = registry_row_to_cursor_server(row)
return {"mcpServers": servers}
def config_payload_to_upsert_fields(payload: dict[str, Any], *, existing: dict[str, Any]) -> dict[str, Any]:
"""Merge Admin edit payload (Cursor-ish) onto an existing registry row."""
base = dict(existing)
sid = str(payload.get("server_id") or base.get("server_id") or "").strip()
if not sid:
raise ValueError("server_id_required")
# Prefer a nested cursor server object when provided.
cursor_server = payload.get("server") if isinstance(payload.get("server"), dict) else None
if cursor_server is None and (
"command" in payload or "url" in payload or "args" in payload or "env" in payload or "headers" in payload
):
cursor_server = {
k: payload.get(k)
for k in ("command", "args", "env", "url", "baseUrl", "headers", "type", "isActive", "enabled", "timeout_s")
if k in payload
}
if cursor_server is not None:
converted = cursor_server_to_install_payload(sid, cursor_server)
if "enabled" in payload:
converted["enabled"] = bool(payload.get("enabled"))
if "timeout_s" in payload:
converted["timeout_s"] = float(payload.get("timeout_s") or 30.0)
return converted
# Field-level updates without full Cursor object.
entry_command = str(payload.get("entry_command") if "entry_command" in payload else base.get("entry_command") or "").strip()
if "entry_args" in payload:
entry_args = [str(x) for x in (payload.get("entry_args") or [])] if isinstance(payload.get("entry_args"), list) else []
else:
raw = base.get("entry_args")
entry_args = [str(x) for x in raw] if isinstance(raw, list) else []
if "env" in payload and isinstance(payload.get("env"), dict):
env_schema = _env_schema_from_cursor_env(payload.get("env"))
elif "env_schema" in payload and isinstance(payload.get("env_schema"), dict):
env_schema = dict(payload.get("env_schema") or {})
else:
env_schema = dict(base.get("env_schema") or {}) if isinstance(base.get("env_schema"), dict) else {}
enabled = bool(payload.get("enabled")) if "enabled" in payload else bool(base.get("enabled"))
timeout_s = float(payload.get("timeout_s") if "timeout_s" in payload else (base.get("timeout_s") or 30.0))
source_type = str(base.get("source_type") or "local")
source_ref = str(base.get("source_ref") or sid)
# If editing looks like mcp-remote, keep npm/mcp-remote markers.
if "mcp-remote" in entry_args or source_ref == "mcp-remote":
source_type = "npm"
source_ref = "mcp-remote"
elif source_type not in {"github", "npm", "pypi", "local"}:
source_type = "local"
source_ref = sid
return {
"server_id": sid,
"source_type": source_type,
"source_ref": source_ref,
"version": str(base.get("version") or ""),
"entry_command": entry_command,
"entry_args": entry_args,
"env_schema": env_schema,
"required_permissions": list(base.get("required_permissions") or [])
if isinstance(base.get("required_permissions"), list)
else [],
"risk_level": str(base.get("risk_level") or "high"),
"enabled": enabled,
"timeout_s": timeout_s,
}
def dumps_cursor_export(rows: list[dict[str, Any]]) -> str:
return json.dumps(build_cursor_mcp_export(rows), ensure_ascii=False, indent=2) + "\n"
__all__ = [
"build_cursor_mcp_export",
"config_payload_to_upsert_fields",
"cursor_server_to_install_payload",
"dumps_cursor_export",
"parse_cursor_mcp_document",
"registry_row_to_cursor_server",
]

View file

@ -1,153 +0,0 @@
from __future__ import annotations
from typing import Any
import time
import httpx
_TRENDING_CACHE: dict[str, Any] = {"ts": 0.0, "items": []}
_TRENDING_TTL_S = 1800
def _safe_get_json(url: str, *, params: dict[str, Any] | None = None, headers: dict[str, str] | None = None) -> dict[str, Any]:
try:
with httpx.Client(timeout=8.0, follow_redirects=True) as c:
r = c.get(url, params=params or {}, headers=headers or {})
if r.status_code != 200:
return {}
obj = r.json()
return obj if isinstance(obj, dict) else {}
except Exception:
return {}
def search_github_repos(query: str, *, limit: int = 8) -> list[dict[str, Any]]:
q = str(query or "").strip()
if not q:
return []
blob = _safe_get_json(
"https://api.github.com/search/repositories",
params={"q": f"{q} mcp server", "sort": "stars", "order": "desc", "per_page": max(1, min(limit, 20))},
headers={"Accept": "application/vnd.github+json"},
)
items = blob.get("items") if isinstance(blob.get("items"), list) else []
out: list[dict[str, Any]] = []
for it in items[:limit]:
if not isinstance(it, dict):
continue
out.append(
{
"source_type": "github",
"name": str(it.get("full_name") or ""),
"source_ref": str(it.get("clone_url") or it.get("html_url") or ""),
"description": str(it.get("description") or ""),
"version": "",
"homepage": str(it.get("html_url") or ""),
"stars": int(it.get("stargazers_count") or 0),
"install_template": infer_install_template("github", str(it.get("clone_url") or it.get("html_url") or "")),
}
)
return out
def search_npm_packages(query: str, *, limit: int = 8) -> list[dict[str, Any]]:
q = str(query or "").strip()
if not q:
return []
blob = _safe_get_json(
"https://registry.npmjs.org/-/v1/search",
params={"text": f"{q} mcp", "size": max(1, min(limit, 20))},
)
items = blob.get("objects") if isinstance(blob.get("objects"), list) else []
out: list[dict[str, Any]] = []
for it in items[:limit]:
pkg = it.get("package") if isinstance(it, dict) else None
if not isinstance(pkg, dict):
continue
out.append(
{
"source_type": "npm",
"name": str(pkg.get("name") or ""),
"source_ref": str(pkg.get("name") or ""),
"description": str(pkg.get("description") or ""),
"version": str(pkg.get("version") or ""),
"homepage": str(pkg.get("links", {}).get("npm") if isinstance(pkg.get("links"), dict) else ""),
"stars": 0,
"install_template": infer_install_template("npm", str(pkg.get("name") or "")),
}
)
return out
def search_pypi_packages(query: str, *, limit: int = 8) -> list[dict[str, Any]]:
q = str(query or "").strip()
if not q:
return []
blob = _safe_get_json(
"https://pypi.org/search/",
params={"q": f"{q} mcp"},
headers={"Accept": "application/json"},
)
# PyPI JSON search API is not officially stable; keep best-effort.
projects = blob.get("projects") if isinstance(blob.get("projects"), list) else []
out: list[dict[str, Any]] = []
for it in projects[:limit]:
if not isinstance(it, dict):
continue
name = str(it.get("name") or "")
out.append(
{
"source_type": "pypi",
"name": name,
"source_ref": name,
"description": str(it.get("description") or ""),
"version": str(it.get("version") or ""),
"homepage": f"https://pypi.org/project/{name}/" if name else "",
"stars": 0,
"install_template": infer_install_template("pypi", name),
}
)
return out
def search_mcp_market(query: str, *, per_source_limit: int = 6) -> list[dict[str, Any]]:
lim = max(1, min(int(per_source_limit or 6), 20))
out: list[dict[str, Any]] = []
out.extend(search_github_repos(query, limit=lim))
out.extend(search_npm_packages(query, limit=lim))
out.extend(search_pypi_packages(query, limit=lim))
return out
def infer_install_template(source_type: str, source_ref: str) -> dict[str, Any]:
st = str(source_type or "").strip().lower()
sr = str(source_ref or "").strip()
if st == "npm":
pkg = sr.split("/")[-1] if sr else ""
return {"entry_command": "npx", "entry_args": [pkg] if pkg else []}
if st == "pypi":
pkg = sr.replace("-", "_")
return {"entry_command": "python", "entry_args": ["-m", pkg] if pkg else []}
return {"entry_command": "python", "entry_args": []}
def trending_mcp_market(*, force_refresh: bool = False, per_source_limit: int = 5) -> list[dict[str, Any]]:
now = time.time()
if not force_refresh and _TRENDING_CACHE["items"] and (now - float(_TRENDING_CACHE["ts"] or 0.0) < _TRENDING_TTL_S):
return list(_TRENDING_CACHE["items"])
items = search_mcp_market("mcp", per_source_limit=per_source_limit)
items = sorted(items, key=lambda x: int(x.get("stars") or 0), reverse=True)
_TRENDING_CACHE["ts"] = now
_TRENDING_CACHE["items"] = list(items)
return items
__all__ = [
"search_mcp_market",
"search_github_repos",
"search_npm_packages",
"search_pypi_packages",
"infer_install_template",
"trending_mcp_market",
]