9.4 Harness 中的 MCP 集成模式
本节讨论在 Harness 系统级别集成多个 MCP Server 时的核心架构、动态工具发现与注册、Schema 缓存策略、权限与审计管理,以及错误处理与降级方案。这些模式确保大规模智能体系统能够高效、安全地管理分布式的 MCP Server。
示例的协议版本:本节的注册中心与编排代码按当前修订版 2026-07-28 编写:没有握手,注册阶段用
server/discover一次拿到能力与supportedVersions,Schema 缓存以结果里的ttlMs/cacheScope为准。当前修订版 2026-07-28 取消了握手(改为每请求_meta自带元数据),服务端也不再主动向客户端发起请求,改用 MRTR:服务端返回input_required结果,客户端补齐输入后用新的 JSON-RPC id 重试原请求。对照见 9.1 协议设计。
9.4.1 系统级集成的核心挑战
当 Harness 需要集成多个 MCP Server 时,面临以下挑战:
动态发现与注册:新的 Server 如何自动被 Harness 发现?
Schema 缓存:如何避免每次都重新获取 Schema(省去 Token 和延迟)?
权限与隔离:不同 Agent 应该访问哪些 Server?
错误处理与降级:某个 Server 故障时如何继续工作?
审计与日志:所有 Tool 调用应该被记录用于审计
输入回传:Server 需要用户输入或模型采样时,Harness 如何把结果送回去?
9.4.2 动态工具注册与发现
MCPToolRegistry
工具注册中心的实现代码如下:
from typing import Callable, Dict, List, Optional, Tuple
from dataclasses import dataclass, field
from enum import Enum
import asyncio
import hashlib
import json
from datetime import datetime, timedelta
PROTOCOL_VERSION = "2026-07-28"
# 无状态请求信封:没有握手,每个请求的 params._meta 自带这些字段
META_PROTOCOL_VERSION = "io.modelcontextprotocol/protocolVersion" # 必填
META_CLIENT_CAPABILITIES = "io.modelcontextprotocol/clientCapabilities" # 必填
META_CLIENT_INFO = "io.modelcontextprotocol/clientInfo" # 可选
META_SERVER_INFO = "io.modelcontextprotocol/serverInfo" # 服务器在结果的_meta中回传
# MRTR: inputRequest 的 method -> 它要求客户端预先声明的能力
INPUT_REQUEST_CAPABILITIES = {
"elicitation/create": {"elicitation": {"form": {}}},
"sampling/createMessage": {"sampling": {}}, # 已进入弃用期,仅为对接旧Server保留
"roots/list": {"roots": {}}, # 同上
}
MISSING_CLIENT_CAPABILITY = -32021 # MissingRequiredClientCapabilityError
DEFAULT_INPUT_REQUIRED_MAX_ROUNDS = 10 # 与官方SDK的默认轮数上限一致
@dataclass
class MCPServerConfig:
"""MCP Server配置"""
server_id: str
server_name: str
transport_type: str # "stdio" | "streamable_http"
endpoint: str # 路径或URL
enabled: bool = True
priority: int = 0 # 优先级(用于多个Server提供相同工具时)
timeout_seconds: int = 30
max_retries: int = 2
tags: List[str] = field(default_factory=list) # 标签化分类
@dataclass
class ServerProfile:
"""一次 server/discover 换来的Server全貌"""
server_id: str
supported_versions: List[str]
capabilities: Dict
instructions: Optional[str] = None
server_info: Dict = field(default_factory=dict)
ttl_ms: Optional[int] = None
cache_scope: str = "private"
discovered_at: datetime = field(default_factory=datetime.now)
def supports(self, version: str) -> bool:
return version in self.supported_versions
@dataclass
class ToolSchema:
"""缓存的工具Schema"""
server_id: str
tool_name: str
description: str
input_schema: Dict
cached_at: datetime
schema_hash: str
ttl_ms: Optional[int] = None # 服务器在列表结果中给出的新鲜度提示
cache_scope: str = "private" # "public" | "private"
auth_context: Optional[str] = None # private结果所属的授权上下文
def is_fresh(self, default_ttl_seconds: int) -> bool:
"""优先采用服务器给出的ttlMs,没有则退回本地默认TTL"""
ttl = (
timedelta(milliseconds=self.ttl_ms)
if self.ttl_ms is not None
else timedelta(seconds=default_ttl_seconds)
)
return datetime.now() - self.cached_at < ttl
def is_reusable_by(self, auth_context: Optional[str], default_ttl_seconds: int) -> bool:
"""private结果只能在同一授权上下文内复用,public结果才可以跨上下文共享"""
if not self.is_fresh(default_ttl_seconds):
return False
if self.cache_scope == "public":
return True
return self.auth_context is not None and self.auth_context == auth_context
class MCPToolRegistry:
"""MCP工具注册中心"""
def __init__(self, client_info: Optional[Dict] = None):
self.servers: Dict[str, MCPServerConfig] = {}
self.tool_cache: Dict[str, ToolSchema] = {}
self.server_clients: Dict[str, object] = {}
self.server_profiles: Dict[str, ServerProfile] = {} # server/discover的结果
self.cache_ttl_seconds = 3600 # 服务器未给出ttlMs时的兜底TTL
self.lock = asyncio.Lock()
self.permission_config: Dict[str, Dict[str, List[str]]] = {} # agent_id -> server_id -> [tool_names]
self.input_handlers: Dict[str, Callable] = {} # MRTR: inputRequest的method -> 输入处理函数
self.client_capabilities: Dict = {} # 随 input_handlers 一并声明
self.client_info = client_info or {"name": "harness-registry", "version": "0.1.0"}
self.max_input_rounds = DEFAULT_INPUT_REQUIRED_MAX_ROUNDS # input_required的轮数上限
def register_input_handler(self, method: str, handler: Callable) -> None:
"""注册MRTR输入处理函数,并同步声明对应的客户端能力
服务器不得为客户端未声明的能力发起inputRequest;反过来,若它需要一项
未声明的能力,会返回-32021并在data.requiredCapabilities中列出所需能力。
"""
if method not in INPUT_REQUEST_CAPABILITIES:
raise ValueError(f"Unsupported inputRequest method: {method}")
self.input_handlers[method] = handler
self.client_capabilities.update(INPUT_REQUEST_CAPABILITIES[method])
async def register_server(self, config: MCPServerConfig) -> Optional[ServerProfile]:
"""注册MCP Server,并用一次 server/discover 探明能力与协议版本"""
async with self.lock:
self.servers[config.server_id] = config
print(f"[Registry] Registered MCP Server: {config.server_name}")
profile = await self._discover_server(config.server_id)
if profile and not profile.supports(PROTOCOL_VERSION):
# 版本谈不拢就先停用,免得后面每个请求都换来一个 -32022
config.enabled = False
print(f"[Registry] {config.server_id} lacks {PROTOCOL_VERSION}: {profile.supported_versions}")
return profile
async def unregister_server(self, server_id: str) -> None:
"""注销MCP Server"""
async with self.lock:
if server_id in self.servers:
del self.servers[server_id]
self.server_profiles.pop(server_id, None)
if server_id in self.server_clients:
# 关闭连接
client = self.server_clients[server_id]
if hasattr(client, 'close'):
await client.close()
del self.server_clients[server_id]
# 清除缓存
self.tool_cache = {
k: v for k, v in self.tool_cache.items()
if v.server_id != server_id
}
async def discover_tools(self, auth_context: Optional[str] = None) -> Dict[str, List[str]]:
"""发现所有可用的工具"""
tools_by_server = {}
for server_id, config in self.servers.items():
if not config.enabled:
continue
profile = self.server_profiles.get(server_id)
if profile and "tools" not in profile.capabilities:
continue # server/discover已经说明它不提供工具,省下一次 tools/list
try:
client = await self._get_client(server_id)
response = await self._send(client, "tools/list")
schemas = await self._cache_tools(server_id, response.get("result", {}), auth_context)
tools_by_server[server_id] = list(schemas)
except Exception as e:
print(f"[Registry] Error discovering tools from {server_id}: {e}")
tools_by_server[server_id] = []
return tools_by_server
async def get_tool_schema(
self,
tool_name: str,
server_id: Optional[str] = None,
auth_context: Optional[str] = None,
) -> Optional[ToolSchema]:
"""获取工具Schema(支持缓存)"""
# 尝试从缓存获取:先看当前授权上下文的私有副本,再看公共副本
async with self.lock:
for sid in ([server_id] if server_id else list(self.servers)):
for scope in ("private", "public"):
cached = self.tool_cache.get(self._cache_key(sid, tool_name, auth_context, scope))
if cached and cached.is_reusable_by(auth_context, self.cache_ttl_seconds):
return cached
# 缓存未命中,从Server获取
if server_id:
servers_to_try = [server_id]
else:
# 尝试所有提供此工具的Server
servers_to_try = []
for sid, config in self.servers.items():
if config.enabled:
servers_to_try.append(sid)
for sid in servers_to_try:
try:
client = await self._get_client(sid)
response = await self._send(client, "tools/list")
schemas = await self._cache_tools(sid, response.get("result", {}), auth_context)
if tool_name in schemas:
return schemas[tool_name]
except Exception as e:
print(f"[Registry] Error getting schema from {sid}: {e}")
continue
return None
async def call_tool(
self,
tool_name: str,
arguments: Dict,
agent_id: Optional[str] = None,
server_id: Optional[str] = None,
) -> Tuple[bool, any]:
"""调用工具,并按MRTR补齐服务器索要的输入"""
try:
# 确定使用哪个Server
if not server_id:
server_id = await self._find_server_for_tool(tool_name, agent_id)
if not server_id:
return False, f"Tool {tool_name} not found in any server"
# 检查权限
if not await self._check_permission(agent_id, server_id, tool_name):
return False, f"Agent {agent_id} not authorized to call {tool_name}"
# 获取client并调用
client = await self._get_client(server_id)
params = {"name": tool_name, "arguments": arguments}
rounds = 0
while True:
# 每次 send_request 都自增JSON-RPC id,所以重试天然用的是新id
response = await self._send(client, "tools/call", params)
if "error" in response:
error = response["error"]
if error.get("code") == MISSING_CLIENT_CAPABILITY:
required = (error.get("data") or {}).get("requiredCapabilities")
return False, f"Client capability not declared: {required} (retry won't help)"
return False, error["message"]
result = response.get("result") or {}
# 老Server不带resultType,按 "complete" 处理
if result.get("resultType", "complete") != "input_required":
return True, result
rounds += 1
if rounds > self.max_input_rounds:
return False, f"input_required exceeded {self.max_input_rounds} rounds"
inputs = await self._collect_inputs(result.get("inputRequests") or {}, agent_id)
if inputs is None:
return False, "No handler for the requested input"
params = {
"name": tool_name,
"arguments": arguments,
"inputResponses": inputs, # 键必须与inputRequests的键一致
"requestState": result.get("requestState"), # 不透明,原样回传
}
except Exception as e:
return False, str(e)
async def _collect_inputs(self, input_requests: Dict, agent_id: Optional[str]) -> Optional[Dict]:
"""按 inputRequests 逐项收集输入(elicitation走审批队列,sampling代调模型)"""
responses = {}
for key, request in input_requests.items():
handler = self.input_handlers.get(request.get("method"))
if handler is None:
print(f"[MRTR] No handler for {request.get('method')} (key={key})")
return None
answer = handler(request.get("params") or {}, agent_id)
if asyncio.iscoroutine(answer):
answer = await answer
if answer is None:
return None
responses[key] = answer
return responses
async def _send(self, client, method: str, params: Optional[Dict] = None) -> Dict:
"""发请求,顺手把无状态信封塞进 params._meta
传输层必须把同一个版本号写进 MCP-Protocol-Version 头,两处不一致的请求会被拒。
"""
payload = dict(params or {})
payload["_meta"] = {
META_PROTOCOL_VERSION: PROTOCOL_VERSION,
META_CLIENT_CAPABILITIES: self.client_capabilities,
META_CLIENT_INFO: self.client_info,
}
# stdio 客户端的 send_request 是同步方法,Streamable HTTP 的是协程
response = client.send_request(method, payload)
if asyncio.iscoroutine(response):
response = await response
return response
async def _discover_server(self, server_id: str) -> Optional[ServerProfile]:
"""server/discover:一次往返拿到 supportedVersions、capabilities 与缓存提示"""
try:
client = await self._get_client(server_id)
response = await self._send(client, "server/discover")
result = response.get("result", {})
profile = ServerProfile(
server_id=server_id,
supported_versions=result.get("supportedVersions", []),
capabilities=result.get("capabilities", {}),
instructions=result.get("instructions"),
server_info=(result.get("_meta") or {}).get(META_SERVER_INFO, {}),
ttl_ms=result.get("ttlMs"),
cache_scope=result.get("cacheScope", "private"),
)
async with self.lock:
self.server_profiles[server_id] = profile
return profile
except Exception as e:
print(f"[Registry] Error discovering {server_id}: {e}")
return None
async def _cache_tools(
self,
server_id: str,
result: Dict,
auth_context: Optional[str],
) -> Dict[str, ToolSchema]:
"""按结果里的 ttlMs / cacheScope 缓存 tools/list,返回本次拿到的全部Schema"""
ttl_ms = result.get("ttlMs")
cache_scope = result.get("cacheScope", "private")
# 私有结果没有授权上下文可绑定时只用一次,不进缓存
cacheable = cache_scope == "public" or bool(auth_context)
schemas = {}
for tool in result.get("tools", []):
schemas[tool["name"]] = ToolSchema(
server_id=server_id,
tool_name=tool["name"],
description=tool.get("description", ""),
input_schema=tool["inputSchema"],
cached_at=datetime.now(),
schema_hash=self._hash_schema(tool),
ttl_ms=ttl_ms,
cache_scope=cache_scope,
auth_context=auth_context if cache_scope == "private" else None,
)
if cacheable:
async with self.lock:
for name, schema in schemas.items():
self.tool_cache[self._cache_key(server_id, name, auth_context, cache_scope)] = schema
return schemas
def _cache_key(self, server_id: str, tool_name: str, auth_context: Optional[str], scope: str) -> str:
"""public结果共用一个键;private结果按授权上下文各存一份,不允许跨上下文复用"""
if scope == "private" and auth_context:
return f"{server_id}#{tool_name}@{auth_context}"
return f"{server_id}#{tool_name}"
async def _get_client(self, server_id: str):
"""获取或创建Server客户端"""
if server_id in self.server_clients:
return self.server_clients[server_id]
config = self.servers.get(server_id)
if not config:
raise ValueError(f"Server {server_id} not found")
if config.transport_type == "stdio":
from mcp_client import StdioMCPClient
client = StdioMCPClient(config.endpoint)
client.start()
elif config.transport_type == "streamable_http":
from mcp_client import StreamableHttpMCPClient
client = StreamableHttpMCPClient(config.endpoint)
await client.connect()
else:
raise ValueError(f"Unknown transport type: {config.transport_type}")
self.server_clients[server_id] = client
return client
async def _find_server_for_tool(self, tool_name: str, auth_context: Optional[str] = None) -> Optional[str]:
"""找到提供某个工具的Server"""
tools_by_server = await self.discover_tools(auth_context)
# 按优先级排序
candidates = [
(sid, self.servers[sid].priority)
for sid, tools in tools_by_server.items()
if tool_name in tools
]
if candidates:
candidates.sort(key=lambda x: x[1], reverse=True)
return candidates[0][0]
return None
async def _check_permission(self, agent_id: Optional[str], server_id: str, tool_name: str) -> bool:
"""检查Agent是否有权限调用工具"""
if agent_id is None:
return False # 匿名调用默认拒绝
# 查询权限配置
allowed_tools = self.permission_config.get(agent_id, {}).get(server_id, [])
# 支持通配符
if "*" in allowed_tools or tool_name in allowed_tools:
return True
print(f"[Permission Denied] agent={agent_id}, server={server_id}, tool={tool_name}")
return False
def _hash_schema(self, tool: Dict) -> str:
"""计算Schema的哈希值,用于判断是否变化"""
schema_str = json.dumps(tool["inputSchema"], sort_keys=True)
return hashlib.md5(schema_str.encode()).hexdigest()
def get_cache_stats(self) -> Dict:
"""获取缓存统计信息"""
return {
"total_cached_tools": len(self.tool_cache),
"registered_servers": len(self.servers),
"profiled_servers": len(self.server_profiles),
"active_clients": len(self.server_clients),
"cache_memory_bytes": sum(len(json.dumps(v.input_schema)) for v in self.tool_cache.values()),
}无握手调用与 MRTR 输入回传
每个请求自带元数据。 上面的注册中心代码已按 2026-07-28 修订版重写,不再有任何握手步骤。这一修订版取消了生命周期握手:initialize 请求与 notifications/initialized 通知都被移除,服务器也不得从同一连接上的历史请求推断状态——需要跨请求保持的状态必须是客户端每次显式传回的标识符,一个长期存活的 stdio 进程并不构成会话。对接这一修订版的 Server 时,客户端不再有独立的初始化步骤,改为在每个请求的 params._meta 中带上必填的 io.modelcontextprotocol/protocolVersion 与 io.modelcontextprotocol/clientCapabilities,并建议带上 io.modelcontextprotocol/clientInfo;需要临时调整日志级别时用 io.modelcontextprotocol/logLevel。缺少必填字段的请求会得到 JSON-RPC -32602(HTTP 400)。服务器信息则建议由服务器放在每个结果 _meta 的 io.modelcontextprotocol/serverInfo 中返回。
现实中两个版本会长期并存,所以对接旧版 Server 的兼容分支短期内不宜删掉,只需要给它一个判定依据:先按新版格式发一个请求,若收到 HTTP 400 就检查响应体——能识别的新版 JSON-RPC 错误说明对方正是新版 Server(纠正后重试即可),响应体为空或无法识别则回退到 initialize 握手。
如果只是想一次性了解某个 Server 的全貌,可以调用 server/discover(服务器必须实现,客户端可选调用):一次往返即可拿到 supportedVersions、capabilities、可选的 instructions、ttlMs、cacheScope 以及 serverInfo,省去 tools/list + prompts/list + resources/list 的多次探测。另外,tools/list、prompts/list、resources/list、resources/read、resources/templates/list 的结果都必须携带 ttlMs(新鲜度提示,单位毫秒)与 cacheScope("public" 或 "private"),这正好可以驱动下面的缓存过期策略:把服务器给出的 ttlMs 当作缓存 TTL,而跨用户共享的缓存层(如远程 L3)只存放 cacheScope 为 "public" 的结果。
服务器不能反向发起请求。 新修订版禁止服务器发送服务器发起的 JSON-RPC 请求,原先的 roots/list、sampling/createMessage、elicitation/create 改走 MRTR 模式:服务器返回 resultType: "input_required" 的结果,其中 inputRequests 是一个映射(服务器分配的键 → 一个 ElicitRequest、CreateMessageRequest 或 ListRootsRequest 请求对象),并带上对客户端不透明的 requestState;客户端收集完输入后,用相同的键构造 inputResponses,连同原样回传的 requestState 重试同一个请求。工程上有几条硬规则:
重试请求的 JSON-RPC
id必须与首次请求不同。客户端不得检查、解析或修改
requestState,只能原样回传。服务器必须至少给出
inputRequests与requestState之一。服务器不得为客户端未声明的能力发起
inputRequest;反过来,若服务器需要一项客户端没有声明的能力,会返回MissingRequiredClientCapabilityError(-32021,HTTP 400),并在data.requiredCapabilities中列出所需能力。这是 Harness 的配置问题(补声明能力,或改派其他 Server),重试无用。只有
prompts/get、resources/read、tools/call可以返回InputRequiredResult。客户端必须把结果中缺失的
resultType当作"complete"(老 Server 不带这个字段)。
对 Harness 来说,这个变化把“人在环”的入口收敛到了同一条路径上:无论是让用户补一个参数(elicitation)、还是让 Harness 替 Server 调一次模型(sampling),都表现为“工具调用返回 input_required,Harness 收集输入后重试”,可以直接复用后面 PermissionGateway 的审批队列,而不必为服务器发起的回调单独维护一条反向通道。同时要给重试设上限:input_required 与重试构成一个循环,没有轮数上限的实现会被一个坏 Server 拖死。
安全提示:requestState 是攻击面。 如果 Harness 自己也写 MCP Server,必须把 requestState 当作攻击者可控的输入:一旦它会影响授权判定、资源访问或业务逻辑,就必须做完整性保护(HMAC 或 AEAD),校验失败直接拒绝;并且应当把调用主体、较短的有效期和发起请求的标识绑定进去,以防重放。客户端一侧则相反——只做透明转发,不做任何解析。
需要说明的是,Roots、Sampling、Logging 只是进入弃用期(至少 12 个月内仍然可用,但新实现不建议采用),并没有被移除。迁移方向是:目录与文件通过工具参数或资源 URI 传入,取代 Roots;直接调用模型提供方 API,取代 Sampling;日志写 stderr 或 OpenTelemetry,取代 Logging。因此新写的输入回传逻辑主要服务于 elicitation,Roots 与 Sampling 两类 inputRequest 只在对接尚未迁移的旧 Server 时才需要处理。
Schema 缓存策略
缓存的多层设计
权限与审计网关
权限与审计网关的实现代码如下:
错误处理与降级策略
错误处理与降级策略的实现代码如下:
本小节小结
Harness 级别的 MCP 集成需要考虑:
动态发现:通过 MCPToolRegistry 自动发现和注册 Server
Schema 缓存:多层缓存设计(内存、磁盘、远程),显著降低延迟和 Token 消耗
权限隔离:PermissionGateway 确保 Agent 只能访问授权的工具
审计追踪:所有 Tool 调用都被记录用于合规性和调试
错误降级:Server 故障时支持后备方案
Schema 缓存可以减少重复工具描述带来的 Token 消耗;具体节省比例取决于工具数量、Schema 大小和缓存命中率,对于大规模智能体系统尤其重要。
下一节将在 MiniHarness 中实现完整的 MCP 客户端集成。
最后更新于
