For the complete documentation index, see llms.txt. This page is also available as Markdown.

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 时,面临以下挑战:

  1. 动态发现与注册:新的 Server 如何自动被 Harness 发现?

  2. Schema 缓存:如何避免每次都重新获取 Schema(省去 Token 和延迟)?

  3. 权限与隔离:不同 Agent 应该访问哪些 Server?

  4. 错误处理与降级:某个 Server 故障时如何继续工作?

  5. 审计与日志:所有 Tool 调用应该被记录用于审计

  6. 输入回传: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/protocolVersionio.modelcontextprotocol/clientCapabilities,并建议带上 io.modelcontextprotocol/clientInfo;需要临时调整日志级别时用 io.modelcontextprotocol/logLevel。缺少必填字段的请求会得到 JSON-RPC -32602(HTTP 400)。服务器信息则建议由服务器放在每个结果 _metaio.modelcontextprotocol/serverInfo 中返回。

现实中两个版本会长期并存,所以对接旧版 Server 的兼容分支短期内不宜删掉,只需要给它一个判定依据:先按新版格式发一个请求,若收到 HTTP 400 就检查响应体——能识别的新版 JSON-RPC 错误说明对方正是新版 Server(纠正后重试即可),响应体为空或无法识别则回退到 initialize 握手。

如果只是想一次性了解某个 Server 的全貌,可以调用 server/discover(服务器必须实现,客户端可选调用):一次往返即可拿到 supportedVersionscapabilities、可选的 instructionsttlMscacheScope 以及 serverInfo,省去 tools/list + prompts/list + resources/list 的多次探测。另外,tools/listprompts/listresources/listresources/readresources/templates/list 的结果都必须携带 ttlMs(新鲜度提示,单位毫秒)与 cacheScope"public""private"),这正好可以驱动下面的缓存过期策略:把服务器给出的 ttlMs 当作缓存 TTL,而跨用户共享的缓存层(如远程 L3)只存放 cacheScope"public" 的结果。

服务器不能反向发起请求。 新修订版禁止服务器发送服务器发起的 JSON-RPC 请求,原先的 roots/listsampling/createMessageelicitation/create 改走 MRTR 模式:服务器返回 resultType: "input_required" 的结果,其中 inputRequests 是一个映射(服务器分配的键 → 一个 ElicitRequest、CreateMessageRequest 或 ListRootsRequest 请求对象),并带上对客户端不透明的 requestState;客户端收集完输入后,用相同的键构造 inputResponses,连同原样回传的 requestState 重试同一个请求。工程上有几条硬规则:

  • 重试请求的 JSON-RPC id 必须与首次请求不同。

  • 客户端不得检查、解析或修改 requestState,只能原样回传。

  • 服务器必须至少给出 inputRequestsrequestState 之一。

  • 服务器不得为客户端未声明的能力发起 inputRequest;反过来,若服务器需要一项客户端没有声明的能力,会返回 MissingRequiredClientCapabilityError-32021,HTTP 400),并在 data.requiredCapabilities 中列出所需能力。这是 Harness 的配置问题(补声明能力,或改派其他 Server),重试无用。

  • 只有 prompts/getresources/readtools/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 集成需要考虑:

  1. 动态发现:通过 MCPToolRegistry 自动发现和注册 Server

  2. Schema 缓存:多层缓存设计(内存、磁盘、远程),显著降低延迟和 Token 消耗

  3. 权限隔离:PermissionGateway 确保 Agent 只能访问授权的工具

  4. 审计追踪:所有 Tool 调用都被记录用于合规性和调试

  5. 错误降级:Server 故障时支持后备方案

Schema 缓存可以减少重复工具描述带来的 Token 消耗;具体节省比例取决于工具数量、Schema 大小和缓存命中率,对于大规模智能体系统尤其重要。

下一节将在 MiniHarness 中实现完整的 MCP 客户端集成。

最后更新于