aps-agent/server/agent_core/mcp_bus.py

776 lines
30 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# ============================================================
# MCP 插件管理总线(moduleId: core-mcp-bus, 可重生 ✅)
# plan.md §6.8 / §7.2 + 矩阵 75 行:在既有外部算法 Skill 基础上扩展 MCP 语义
# - manifest(plugin_id/name/version/tools/description/system)
# - 工具契约(input/output JSON Schema + 权力等级 P0-P3)
# - tool 级 allow/deny 权限(越权 → 拒绝 + 审计,fail closed)
# - 健康探测、审计留痕、启停、版本兼容
# 兼容性:只增不删——SkillRegistry 仍独立工作,本模块可选桥接其清单为 MCP 插件
# 落盘:桌面 ~/.aps/mcp_plugins/plugins.json;Web server/data/mcp_plugins/plugins.json
# (测试可用 APS_MCP_BUS_PATH 或显式 path 指向临时文件)
# ============================================================
from __future__ import annotations
import json
import logging
import os
import tempfile
import threading
from collections.abc import Callable
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
import httpx
from pydantic import BaseModel, Field
from server.timeutil import fmt_dt
# 总线版本:插件声明 min_bus_version <= BUS_VERSION 才兼容
BUS_VERSION = "1.0"
POWERS = ("P0", "P1", "P2", "P3")
_MCP_BUS_PATH_ENV = "APS_MCP_BUS_PATH"
_HEALTH_TIMEOUT = 5.0
_TOOL_TIMEOUT_DEFAULT = 60.0
_logger = logging.getLogger(__name__)
def _now() -> str:
# DTZ005:统一取本地时区的 aware 时间,墙钟与仓库其余模块一致
return fmt_dt(datetime.now(UTC).astimezone())
def _version_tuple(version: str) -> tuple[int, ...]:
"""'1.2.3' → (1,2,3);非法段按 0 处理,保证可比较。"""
out: list[int] = []
for seg in str(version or "").strip().split("."):
digits = "".join(ch for ch in seg if ch.isdigit())
out.append(int(digits) if digits else 0)
return tuple((out + [0, 0, 0])[:3])
def _bus_compatible(min_bus_version: str) -> bool:
return _version_tuple(min_bus_version) <= _version_tuple(BUS_VERSION)
def _default_allow(power: str) -> bool:
"""缺省权限:P0/P1 只读/计算放行;P2/P3 执行/回写默认拒绝(安全优先)。"""
return power in ("P0", "P1")
class McpToolSpec(BaseModel):
"""MCP 工具契约(plan.md §7.2 McpTool 的运行时形态)。"""
name: str
description: str = ""
power: str = "P1" # P0-P3(执行/回写 P2/P3)
input_schema: dict[str, Any] = Field(default_factory=dict) # JSON Schema
output_schema: dict[str, Any] = Field(default_factory=dict) # JSON Schema
idempotent: bool = False
transport: str = "local" # local | http(http 走 {endpoint}/mcp/tools/{name})
class ToolPermission(BaseModel):
"""tool 级 allow/deny 权限(缺省按权力:P0/P1 放行,P2/P3 拒绝)。"""
tool: str
allow: bool = True
reason: str = ""
updatedAt: str = ""
class McpPluginManifest(BaseModel):
"""MCP 插件清单(plan.md §7.2 McpPluginManifest 的运行时形态)。"""
plugin_id: str
name: str
description: str = ""
version: str = "1.0.0"
system: str = "MCP" # MES/QMS/EMS/TMS/WMS/ERP/ALGO/…
endpoint: str = "local://stub"
auth: str = ""
enabled: bool = True
max_power: str = "P1" # 插件整体权力上限(展示/审计用)
min_bus_version: str = "1.0" # 版本兼容:不高于 BUS_VERSION
tools: list[McpToolSpec] = Field(default_factory=list)
events: list[str] = Field(default_factory=list)
permissions: dict[str, ToolPermission] = Field(default_factory=dict)
def tool(self, tool_name: str) -> McpToolSpec | None:
return next((t for t in self.tools if t.name == tool_name), None)
def permission_for(self, tool_name: str) -> ToolPermission:
"""权限解析:显式配置优先,缺省按工具权力。"""
if tool_name in self.permissions:
return self.permissions[tool_name]
spec = self.tool(tool_name)
power = spec.power if spec else self.max_power
return ToolPermission(tool=tool_name, allow=_default_allow(power))
def normalized(self) -> McpPluginManifest:
"""落盘前把缺省权限显式化(工具级 allow/deny 全量可管理)。"""
perms: dict[str, ToolPermission] = {}
for t in self.tools:
perms[t.name] = (
self.permissions[t.name] if t.name in self.permissions
else ToolPermission(tool=t.name, allow=_default_allow(t.power))
)
self.permissions = perms
return self
class McpBusError(Exception):
"""MCP 总线错误基类。"""
class PluginNotFoundError(McpBusError):
pass
class ToolNotFoundError(McpBusError):
pass
class ToolExecutionError(McpBusError):
pass
def _default_plugins() -> list[dict[str, Any]]:
"""内置演示插件:mcp.stub(local://stub + stub.echo,供联调与黄金测试)。"""
return [McpPluginManifest(
plugin_id="mcp.stub",
name="内置 MCP 桩(演示)",
description="本地桩插件:echo 工具用于联调、管理台演示与黄金测试",
system="MCP",
endpoint="local://stub",
version="1.0.0",
max_power="P1",
tools=[McpToolSpec(
name="stub.echo",
description="回显入参(幂等只读)",
power="P1",
input_schema={"type": "object", "properties": {"message": {"type": "string"}}},
output_schema={"type": "object", "properties": {"echo": {"type": "string"}}},
idempotent=True,
transport="local",
)],
).model_dump(mode="json")]
def _default_bus_path() -> Path:
from server.aps_home import aps_home
return aps_home() / "mcp_plugins" / "plugins.json"
def from_skill_manifest(skill: dict[str, Any]) -> dict[str, Any]:
"""既有外部算法 Skill → MCP 插件清单(向后兼容桥接,只增不删)。"""
sid = skill.get("skill_id") or "unknown"
endpoint = skill.get("endpoint") or "local://stub"
max_power = str(skill.get("max_power") or "P1")
if max_power not in POWERS:
max_power = "P1"
return McpPluginManifest(
plugin_id=f"skill.{sid}",
name=str(skill.get("name") or sid),
description=str(skill.get("description") or ""),
version=str(skill.get("version") or "1.0"),
system="ALGO",
endpoint=endpoint,
auth=str(skill.get("auth") or ""),
enabled=bool(skill.get("enabled", True)),
max_power=max_power,
tools=[McpToolSpec(
name="flex.schedule",
description=f"外部算法排产(桥接自 Skill {sid})",
power=max_power,
transport="local" if endpoint.startswith("local://") else "http",
idempotent=False,
)],
).model_dump(mode="json")
# ---------------- 真实 MES HTTP 适配器(矩阵 75:现场对接只需配置) ----------------
def mes_http_adapter_manifest(base_url: str = "") -> dict[str, Any]:
"""真实 MES 适配器 manifest(plugin_id: mes.http,工具契约 + power 门禁)。
- 未配置 base_url 时 endpoint 保持 local://http-adapter,工具调用 fail-closed;
- token 不入库(auth 留空),由 HttpMesClient 按环境变量读取。
"""
return McpPluginManifest(
plugin_id="mes.http",
name="\u771f\u5b9e MES HTTP \u9002\u914d\u5668",
description="config-driven HTTP MES \u9002\u914d\u5668\uff1a\u9274\u6743/\u5e42\u7b49/\u8d85\u65f6/\u91cd\u8bd5/\u56de\u6267\uff1b"
"\u672a\u914d\u7f6e MES_HTTP_BASE_URL \u65f6\u5de5\u5177\u8c03\u7528 fail-closed\uff08\u660e\u786e\u62a5\u9519\uff09\u3002",
system="MES",
endpoint=base_url or "local://http-adapter",
auth="",
enabled=True,
max_power="P3",
min_bus_version="1.0",
tools=[
McpToolSpec(
name="mes.http_dispatch",
description="\u771f\u5b9e MES \u4e0b\u53d1\uff08\u5e42\u7b49 idemKey\uff0cP3 \u95e8\u7981\uff09",
power="P3",
input_schema={
"type": "object",
"properties": {
"idemKey": {"type": "string", "description": "\u5e42\u7b49\u952e"},
"track": {"type": "string"},
"apsWoId": {"type": "string"},
"orderNo": {"type": "string"},
"operation": {"type": "string"},
"equipment": {"type": "string"},
"start": {"type": "string"},
"end": {"type": "string"},
"versionId": {"type": "integer"},
"versionNo": {"type": "string"},
},
"required": ["idemKey"],
},
output_schema={"type": "object"},
idempotent=True,
transport="local",
),
McpToolSpec(
name="mes.http_status",
description="\u771f\u5b9e MES \u8fde\u63a5\u4e0e\u5916\u90e8\u5de5\u5355\u72b6\u6001\u67e5\u8be2\uff08\u53ea\u8bfb\uff09",
power="P0",
input_schema={
"type": "object",
"properties": {"externalWoId": {"type": "string"}},
},
output_schema={"type": "object"},
idempotent=True,
transport="local",
),
McpToolSpec(
name="mes.http_report",
description="\u771f\u5b9e MES \u62a5\u5de5\u56de\u6267\uff08\u8fdb\u5ea6/\u6570\u91cf/\u72b6\u6001\u56de\u6d41\uff09",
power="P2",
input_schema={
"type": "object",
"properties": {
"externalWoId": {"type": "string"},
"progressPct": {"type": "integer"},
"qtyDone": {"type": "number"},
"status": {"type": "string"},
},
"required": ["externalWoId"],
},
output_schema={"type": "object"},
idempotent=False,
transport="local",
),
McpToolSpec(
name="mes.http_cancel",
description="\u771f\u5b9e MES \u64a4\u5355\uff08\u5e42\u7b49\uff0csaga \u8865\u507f\uff09",
power="P2",
input_schema={
"type": "object",
"properties": {
"externalWoId": {"type": "string"},
"reason": {"type": "string"},
},
"required": ["externalWoId"],
},
output_schema={"type": "object"},
idempotent=True,
transport="local",
),
],
).model_dump(mode="json")
def _mes_http_tool_handlers() -> dict[str, Callable[[dict[str, Any]], Any]]:
"""HttpMesClient \u5de5\u5177 handler\uff08\u61d2\u89e3\u6790\u8fdb\u7a0b\u5185\u5355\u4f8b\uff0c\u914d\u7f6e\u70ed\u5207\u6362\u53ef reset \u91cd\u5efa\uff09\u3002"""
from server.integrations.mes_http import get_http_mes_client
def _dispatch(args: dict[str, Any]) -> Any:
client = get_http_mes_client()
payload = {k: args.get(k) for k in (
"track", "apsWoId", "orderNo", "operation", "equipment",
"start", "end", "versionId", "versionNo",
) if k in args}
return client.create_work_order(payload, str(args.get("idemKey") or ""))
def _status(args: dict[str, Any]) -> Any:
client = get_http_mes_client()
external_wo_id = args.get("externalWoId")
if not external_wo_id:
return client.status()
return client.fetch_work_order(str(external_wo_id))
def _report(args: dict[str, Any]) -> Any:
client = get_http_mes_client()
external_wo_id = str(args.get("externalWoId") or "")
if not external_wo_id:
raise ValueError("mes.http_report \u7f3a\u5c11 externalWoId")
payload = {k: args[k] for k in ("progressPct", "qtyDone", "status") if k in args}
return client.post_report(external_wo_id, payload)
def _cancel(args: dict[str, Any]) -> Any:
client = get_http_mes_client()
external_wo_id = str(args.get("externalWoId") or "")
if not external_wo_id:
raise ValueError("mes.http_cancel \u7f3a\u5c11 externalWoId")
return client.cancel_work_order(
external_wo_id, reason=str(args.get("reason") or "saga-compensation"))
return {
"mes.http_dispatch": _dispatch,
"mes.http_status": _status,
"mes.http_report": _report,
"mes.http_cancel": _cancel,
}
class McpBus:
"""MCP 插件注册表:manifest / 工具契约 / 权限 / 健康 / 审计 / 启停 / 版本兼容。"""
BUS_VERSION = BUS_VERSION
def __init__(
self,
path: str | None = None,
*,
audit_sink: Callable[[dict[str, Any]], None] | None = None,
) -> None:
self.legacy_path = path or os.environ.get(_MCP_BUS_PATH_ENV)
self.path = Path(self.legacy_path) if self.legacy_path else _default_bus_path()
self.audit_sink = audit_sink
self._lock = threading.RLock()
self.audit_events: list[dict[str, Any]] = []
self.health_history: dict[str, list[dict[str, Any]]] = {}
self.call_stats: dict[str, dict[str, dict[str, Any]]] = {}
self.handlers: dict[tuple[str, str], Callable[[dict[str, Any]], Any]] = {}
self.plugins: list[McpPluginManifest] = self._load()
if not self.plugins:
self.plugins = [McpPluginManifest(**m) for m in _default_plugins()]
self._write_all()
# 内置桩处理器
self.handlers[("mcp.stub", "stub.echo")] = (
lambda args: {"echo": args.get("message", "")}
)
# 真实 MES 适配器(矩阵 75):mes.http 工具注册 + 本地 handler(fail-closed)
self.seed_mes_http_adapter(actor="system")
# ---------------- 真实 MES 适配器(矩阵 75) ----------------
def seed_mes_http_adapter(self, *, actor: str = "system", overwrite: bool = True) -> dict[str, Any]:
"""注册真实 MES HTTP 适配器(mes.http):manifest + power 门禁 + 本地 handler。
未配置 MES_HTTP_BASE_URL 时工具调用 fail-closed(MesHttpError MES_HTTP_NOT_CONFIGURED)。
"""
base_url = os.environ.get("MES_HTTP_BASE_URL") or ""
manifest = mes_http_adapter_manifest(base_url=base_url)
self.register(manifest, actor=actor, overwrite=overwrite)
for tool, fn in _mes_http_tool_handlers().items():
self.register_handler("mes.http", tool, fn)
return self.get("mes.http") or {}
# ---------------- 持久化 ----------------
def _load(self) -> list[McpPluginManifest]:
try:
data = json.loads(self.path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return []
out: list[McpPluginManifest] = []
for item in data.get("plugins") or []:
try:
out.append(McpPluginManifest(**item))
except (TypeError, ValueError):
continue
return out
def _write_all(self) -> None:
self.path.parent.mkdir(parents=True, exist_ok=True)
payload = {
"busVersion": BUS_VERSION,
"plugins": [p.model_dump(mode="json") for p in self.plugins],
}
fd, tmp = tempfile.mkstemp(dir=str(self.path.parent), suffix=".tmp")
try:
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump(payload, f, ensure_ascii=False, indent=1)
os.replace(tmp, self.path)
except BaseException:
if os.path.exists(tmp):
os.unlink(tmp)
raise
# ---------------- 审计(本总线审计 + 可选全局审计链 sink) ----------------
def _audit(
self,
*,
actor: str,
action: str,
target: dict[str, Any],
power: str,
rationale: dict[str, Any],
result: str = "SUCCESS",
) -> dict[str, Any]:
event: dict[str, Any] = {
"at": _now(),
"actor": actor,
"category": "MCP",
"action": action,
"target": target,
"power": power,
"rationale": rationale,
"result": result,
}
self.audit_events.append(event)
if self.audit_sink is not None:
try:
self.audit_sink(event) # 接入 write_audit 证据链(尽力而为)
except Exception:
_logger.warning("MCP 审计 sink 失败:%s", event.get("action"), exc_info=True)
return event
def audit(self, plugin_id: str | None = None, limit: int = 50) -> list[dict[str, Any]]:
rows = list(self.audit_events)
if plugin_id:
rows = [e for e in rows if (e.get("target") or {}).get("id") == plugin_id]
return list(reversed(rows))[: max(1, min(limit, 500))]
# ---------------- 注册表 ----------------
def get_manifest(self, plugin_id: str) -> McpPluginManifest:
for p in self.plugins:
if p.plugin_id == plugin_id:
return p
raise PluginNotFoundError(f"未找到 MCP 插件:{plugin_id}")
def get(self, plugin_id: str) -> dict[str, Any] | None:
try:
p = self.get_manifest(plugin_id)
except PluginNotFoundError:
return None
return p.model_dump(mode="json")
def list(self) -> list[dict[str, Any]]:
return [p.model_dump(mode="json") for p in self.plugins]
def describe(self, plugin_id: str) -> dict[str, Any] | None:
p = self.get(plugin_id)
if p is None:
return None
return {
**p,
"compatible": _bus_compatible(str(p.get("min_bus_version") or "1.0")),
"callStats": self._stats_for(plugin_id),
}
def register(
self,
manifest: dict[str, Any] | McpPluginManifest,
*,
actor: str = "planner",
overwrite: bool = True,
) -> dict[str, Any]:
"""登记/更新插件:版本兼容校验 + 工具契约校验 + 缺省权限落盘 + 审计。
覆盖登记时保留既有工具的自定义权限(只增不删,避免重装丢失授权)。
"""
m = manifest if isinstance(manifest, McpPluginManifest) else McpPluginManifest(**manifest)
if not m.plugin_id or not m.name:
raise ValueError("plugin_id 与 name 必填")
if not m.tools:
raise ValueError(f"MCP 插件 {m.plugin_id} 至少需要一个工具(tools 不能为空)")
seen: set[str] = set()
for t in m.tools:
if not t.name or "/" in t.name:
raise ValueError(f"工具名非法:{t.name!r}")
if t.name in seen:
raise ValueError(f"工具名重复:{t.name}")
seen.add(t.name)
if not _bus_compatible(m.min_bus_version):
raise ValueError(
f"插件 {m.plugin_id} 需要总线版本 ≥ {m.min_bus_version},"
f"当前总线 {BUS_VERSION},拒绝登记(版本不兼容)"
)
with self._lock:
existing = next((p for p in self.plugins if p.plugin_id == m.plugin_id), None)
if existing is not None and not overwrite:
raise ValueError(f"MCP 插件已存在:{m.plugin_id}")
if existing is not None:
# 保留既有工具的自定义权限(新工具走缺省)
for t in m.tools:
if t.name in existing.permissions and t.name not in m.permissions:
m.permissions[t.name] = existing.permissions[t.name]
m = m.normalized()
if existing is None:
self.plugins.append(m)
else:
for i, p in enumerate(self.plugins):
if p.plugin_id == m.plugin_id:
self.plugins[i] = m
break
self._write_all()
self._audit(
actor=actor, action="mcp.plugin.register",
target={"type": "MCP_PLUGIN", "id": m.plugin_id},
power="P2",
rationale={"version": m.version, "tools": [t.name for t in m.tools]},
)
return m.model_dump(mode="json")
def unregister(self, plugin_id: str, *, actor: str = "planner") -> bool:
with self._lock:
before = len(self.plugins)
self.plugins = [p for p in self.plugins if p.plugin_id != plugin_id]
if len(self.plugins) == before:
return False
self._write_all()
self._audit(
actor=actor, action="mcp.plugin.unregister",
target={"type": "MCP_PLUGIN", "id": plugin_id},
power="P2", rationale={},
)
return True
def set_enabled(self, plugin_id: str, enabled: bool, *, actor: str = "planner") -> dict[str, Any]:
"""启停插件:停用后 call_tool 一律拒绝并审计(fail closed)。"""
with self._lock:
p = self.get_manifest(plugin_id)
p.enabled = bool(enabled)
self._write_all()
self._audit(
actor=actor,
action="mcp.plugin.start" if p.enabled else "mcp.plugin.stop",
target={"type": "MCP_PLUGIN", "id": plugin_id},
power="P2", rationale={"enabled": p.enabled},
)
return p.model_dump(mode="json")
def set_permission(
self,
plugin_id: str,
tool: str,
allow: bool,
*,
actor: str = "planner",
reason: str = "",
) -> dict[str, Any]:
"""tool 级权限开关:更新后立即生效,变更写审计。"""
with self._lock:
p = self.get_manifest(plugin_id)
if p.tool(tool) is None:
raise ToolNotFoundError(f"插件 {plugin_id} 无工具 {tool}")
p.permissions[tool] = ToolPermission(
tool=tool, allow=bool(allow), reason=reason or "", updatedAt=_now(),
)
self._write_all()
self._audit(
actor=actor, action="mcp.permission.update",
target={"type": "MCP_PLUGIN", "id": plugin_id, "tool": tool},
power="P2", rationale={"allow": bool(allow), "reason": reason},
)
return p.model_dump(mode="json")
# ---------------- 健康 ----------------
def health(self, plugin_id: str | None = None) -> list[dict[str, Any]]:
if plugin_id is not None:
try:
targets = [self.get_manifest(plugin_id)]
except PluginNotFoundError:
return []
else:
targets = list(self.plugins)
out: list[dict[str, Any]] = []
for p in targets:
status = self._probe(p)
record = {"ts": _now(), **status}
hist = self.health_history.setdefault(p.plugin_id, [])
hist.append(record)
del hist[:-20] # 只留最近 20 次
out.append({
"plugin_id": p.plugin_id,
"name": p.name,
"version": p.version,
"enabled": p.enabled,
"endpoint": p.endpoint,
**status,
})
return out
def history(self, plugin_id: str) -> list[dict[str, Any]]:
return list(self.health_history.get(plugin_id) or [])
def _probe(self, p: McpPluginManifest) -> dict[str, Any]:
ep = p.endpoint or ""
if ep.startswith("local://"):
return {"ok": True, "latencyMs": 0, "detail": "local stub"}
health_url = ep.rstrip("/") + "/health"
try:
headers = {}
if p.auth:
headers["Authorization"] = f"Bearer {p.auth}"
with httpx.Client(timeout=_HEALTH_TIMEOUT) as client:
r = client.get(health_url, headers=headers)
return {
"ok": r.status_code < 400,
"latencyMs": int(r.elapsed.total_seconds() * 1000),
"detail": f"HTTP {r.status_code}",
}
except (httpx.HTTPError, OSError, ValueError) as exc:
return {"ok": False, "latencyMs": None, "detail": str(exc)}
# ---------------- 版本兼容 ----------------
def version_report(self) -> list[dict[str, Any]]:
return [{
"plugin_id": p.plugin_id,
"name": p.name,
"version": p.version,
"min_bus_version": p.min_bus_version,
"bus_version": BUS_VERSION,
"compatible": _bus_compatible(p.min_bus_version),
} for p in self.plugins]
# ---------------- 调用统计 ----------------
def _bump(self, plugin_id: str, tool: str, outcome: str, *, ok: bool | None = None) -> None:
tools = self.call_stats.setdefault(plugin_id, {})
row = tools.setdefault(tool, {"allowed": 0, "denied": 0, "failed": 0, "lastTs": None})
row[outcome] = int(row.get(outcome, 0)) + 1
row["lastTs"] = _now()
if ok is not None:
row["lastOk"] = ok
def _stats_for(self, plugin_id: str) -> dict[str, Any]:
return dict(self.call_stats.get(plugin_id) or {})
def stats(self, plugin_id: str | None = None) -> list[dict[str, Any]]:
out: list[dict[str, Any]] = []
for pid, tools in self.call_stats.items():
if plugin_id and pid != plugin_id:
continue
total = sum(
int(row.get("allowed", 0)) + int(row.get("denied", 0)) + int(row.get("failed", 0))
for row in tools.values()
)
out.append({"plugin_id": pid, "tools": tools, "totalCalls": total})
return out
# ---------------- 工具调用(权限门禁 + 审计) ----------------
def register_handler(
self,
plugin_id: str,
tool: str,
fn: Callable[[dict[str, Any]], Any],
) -> None:
self.handlers[(plugin_id, tool)] = fn
def call_tool(
self,
plugin_id: str,
tool: str,
args: dict[str, Any] | None = None,
*,
actor: str = "planner",
) -> Any:
"""统一工具入口:插件启停 + tool 权限双重门禁,越权拒绝并写 DENIED 审计。"""
with self._lock:
p = self.get_manifest(plugin_id)
spec = p.tool(tool)
if spec is None:
raise ToolNotFoundError(f"插件 {plugin_id} 无工具 {tool}")
permission = p.permission_for(tool)
target = {"type": "MCP_PLUGIN", "id": plugin_id, "tool": tool}
if not p.enabled:
self._bump(plugin_id, tool, "denied", ok=False)
self._audit(
actor=actor, action="mcp.tool.denied", target=target,
power=spec.power, rationale={"reason": "plugin-disabled"},
result="DENIED",
)
raise PermissionError(f"MCP 插件 {plugin_id} 已停用,拒绝调用工具 {tool}")
if not permission.allow:
self._bump(plugin_id, tool, "denied", ok=False)
self._audit(
actor=actor, action="mcp.tool.denied", target=target,
power=spec.power,
rationale={
"reason": "permission-denied",
"permission": permission.model_dump(mode="json"),
},
result="DENIED",
)
raise PermissionError(
f"工具 {plugin_id}.{tool} 未授权(allow=False),已拒绝并审计"
)
self._audit(
actor=actor, action="mcp.tool.run", target=target,
power=spec.power,
rationale={"idempotent": spec.idempotent, "transport": spec.transport},
)
# 门禁通过:锁外执行(长任务不阻塞注册表其他操作)
try:
result = self._execute(p, spec, args or {}, actor=actor)
except Exception as exc:
self._bump(plugin_id, tool, "failed", ok=False)
self._audit(
actor=actor, action="mcp.tool.failed", target=target,
power=spec.power, rationale={"error": str(exc)}, result="FAILED",
)
raise
self._bump(plugin_id, tool, "allowed", ok=True)
return result
def _execute(self, p: McpPluginManifest, spec: McpToolSpec, args: dict[str, Any], actor: str) -> Any:
if spec.transport == "local" or p.endpoint.startswith("local://"):
handler = self.handlers.get((p.plugin_id, spec.name))
if handler is None:
raise ToolExecutionError(f"本地工具 {p.plugin_id}.{spec.name} 未注册处理器")
return handler(args)
url = p.endpoint.rstrip("/") + "/mcp/tools/" + spec.name
headers = {"Content-Type": "application/json"}
if p.auth:
headers["Authorization"] = f"Bearer {p.auth}"
timeout = float(os.environ.get("APS_MCP_TOOL_TIMEOUT", _TOOL_TIMEOUT_DEFAULT))
with httpx.Client(timeout=timeout) as client:
r = client.post(
url,
json={"tool": spec.name, "arguments": args, "actor": actor},
headers=headers,
)
r.raise_for_status()
return r.json()
# ---------------- Skill 桥接(矩阵 75:既有 skills 扩展 MCP 语义) ----------------
def seed_from_skills(
self,
registry: Any,
*,
actor: str = "system",
overwrite: bool = False,
) -> int:
"""把 SkillRegistry 中每个 Skill 桥接为 MCP 插件(skill.<skill_id>,只增不删)。"""
count = 0
for skill in registry.list():
manifest = from_skill_manifest(skill)
if not overwrite and self.get(manifest["plugin_id"]) is not None:
continue
self.register(manifest, actor=actor, overwrite=overwrite)
count += 1
return count
_bus: McpBus | None = None
def get_mcp_bus() -> McpBus:
global _bus
if _bus is None:
_bus = McpBus()
return _bus
def reset_mcp_bus() -> None:
"""测试用:丢弃单例,下次按当前环境变量重建。"""
global _bus
_bus = None