Agent Harness 工程:接入远程 MCP
Agent Harness 工程:接入远程 MCP
上一章已经打通了一条本地 MCP 工具链:
Agent Harness
-> stdio Client
-> Task 子进程
-> Worktree 内的工具
这种方式很适合文件、Shell 和测试工具。每个 Task 拥有独立进程,连接生命周期也自然地跟随 Worktree。
但并不是所有能力都应该随 Task 启动。例如:
- Issue、日历和知识库由团队共享;
- 数据库网关需要集中管理连接池和审计;
- 内部平台已经是独立服务,不应为每个 Agent 再启动一份;
- 同一个 MCP Server 可能同时服务多个用户和多个 Harness 实例。
此时,MCP Server 会移到网络另一端:
Agent Harness
-> Streamable HTTP Client
-> Remote MCP Server
-> shared service
表面上只是把 stdio_client() 换成 streamable_http_client(),实际上有三条边界同时发生了变化:
| 边界 | 本地 stdio | 远程 Streamable HTTP |
|---|---|---|
| 生命周期 | 子进程通常跟随 Task | Server 独立运行,Client 会断线和重连 |
| 身份 | 主要依靠进程、Worktree 与沙箱隔离 | 每个 HTTP 请求都要验证用户、租户和权限 |
| 故障 | 进程退出通常意味着连接结束 | 网络中断不代表请求未执行,也不代表已取消 |
因此,本章不会重写上一章的 Tool Binding、Schema 校验和 Tool Registry,而是为它们补上一条可靠的远程连接路径:
- 理解 Streamable HTTP 中 POST、JSON 与 SSE 的关系;
- 先实现一个无状态远程 Server;
- 使用 OAuth 连接受保护的 MCP Resource Server;
- 将 Token、Session 和 Tenant 放在正确的隔离层;
- 区分流恢复、会话重建与工具重试;
- 为共享连接增加并发、超时和可观测性边界。
Streamable HTTP 不是“把 SSE 地址换个名字”
按照 MCP Streamable HTTP 规范,Server 对外提供一个 MCP endpoint,并同时支持 POST 和 GET。
Client 发送每条 MCP 消息时,都会发起一个新的 HTTP POST:
POST /mcp
Content-Type: application/json
Accept: application/json, text/event-stream
Server 可以根据这次请求的特点返回:
application/json:直接返回一个 JSON-RPC 响应;text/event-stream:在 SSE 流中发送通知、进度和最终响应;202 Accepted:已经接受 notification 或 response,本次没有响应体。
GET /mcp 则可以建立一条由 Server 主动发送消息的 SSE 流。它是可选能力,不是所有 Server 都必须提供。
每条 Client 消息仍由 POST 发出;SSE 是一种可流式返回消息的响应形式,也可以由 GET 建立 Server 主动消息通道。
这里最容易出现三个误解。
误解一:SSE 是另一套工具调用协议
不是。无论 HTTP response 是 JSON 还是 SSE,里面承载的仍然是同一套 MCP JSON-RPC 消息。ClientSession.call_tool() 的调用方式不会因此改变。
误解二:SSE 断开等于工具已取消
也不是。连接断开只说明 Client 暂时收不到后续事件,Server 端的工具可能仍在执行。要取消请求,需要发送 MCP 的取消 notification;直接关闭浏览器、HTTP stream 或 Harness 进程,都不能证明副作用已经停止。
误解三:Streamable HTTP 就是旧的 HTTP+SSE 传输
当前规范中的 Streamable HTTP 已经取代旧的 HTTP+SSE 传输。新实现应连接统一的 MCP endpoint,不要再为旧协议单独拼接一个“POST 消息地址”和一个“SSE 地址”。
先从无状态 Server 开始
共享服务并不一定需要 MCP Session。对于“收到请求、执行工具、立即返回”的服务,最简单的部署方式是:
from typing import Annotated
from mcp.server.fastmcp import FastMCP
from pydantic import BaseModel, Field
mcp = FastMCP(
"issue-query",
stateless_http=True,
json_response=True,
)
class Issue(BaseModel):
id: str
title: str
status: str
ISSUES = {
"ISSUE-42": Issue(
id="ISSUE-42",
title="Add MCP health checks",
status="open",
),
}
@mcp.tool()
async def get_issue(
issue_id: Annotated[
str,
Field(pattern=r"^ISSUE-[0-9]+$"),
],
) -> Issue:
"""Read one issue by its stable ID."""
issue = ISSUES.get(issue_id)
if issue is None:
raise ValueError("issue does not exist")
return issue
if __name__ == "__main__":
mcp.run(transport="streamable-http")
stateless_http=True 表示 Server 不依赖跨请求的 MCP Session;json_response=True 表示普通请求优先返回 JSON,而不是为了一个很短的结果也保持 SSE stream。
这不是功能缩水。无状态模式反而更容易:
- 放在反向代理和负载均衡器后面;
- 横向扩容多个副本;
- 避免进程内 Session 丢失;
- 为每个请求独立验证 Access Token;
- 控制长连接数量。
只有当 Server 确实需要 Server-to-Client notification、长时间进度流、resumability 或跨请求状态时,才应该启用 stateful session。不要先打开状态,再寻找状态存在的理由。
上面的代码适合本机验证协议。生产环境还需要 TLS、认证、Origin 校验、请求大小限制和受控网络出口,不能直接把未认证的开发端口暴露到公网。
Session、Token 与业务身份不是一回事
远程 MCP 中常见的三个标识分别解决不同问题:
| 标识 | 解决的问题 | 能否替代认证 |
|---|---|---|
| Access Token | 当前请求代表谁、被授予哪些权限 | 它就是认证凭据 |
MCP-Session-Id | 多个 HTTP 请求属于哪个协议会话 | 不能 |
| Provider tool call ID | 模型消息中的调用与结果如何闭环 | 不能 |
如果 Server 在初始化响应中返回 MCP-Session-Id,Client 必须在后续请求中带回它。Session ID 应该是不可预测、会过期的 opaque value,但它仍然不是登录凭证。
Server 必须在每一个 HTTP 请求上重新验证 Access Token,并确认该 Session 仍然属于同一用户和同一租户。只在创建 Session 时检查一次 Token,后续仅凭 Session ID 放行,会把 Session 变成一枚隐蔽的长期凭据。
OAuth 决定请求拥有什么权限;MCP Session 只关联协议状态;Server 访问下游系统时还要使用自己的下游凭据,不能转发 Client Token。
MCP Server 是 OAuth Resource Server
按照 MCP 授权规范,在受保护的 HTTP 部署中,MCP Server 扮演 OAuth 2.1 Resource Server,认证服务是独立的 Authorization Server,Harness 则是 OAuth Client。
一个受保护的 FastMCP Server 可以这样装配:
from mcp.server.auth.provider import TokenVerifier
from mcp.server.auth.settings import AuthSettings
from mcp.server.fastmcp import FastMCP
from pydantic import AnyHttpUrl
def build_issue_server(
token_verifier: TokenVerifier,
) -> FastMCP:
server = FastMCP(
"issue-service",
stateless_http=True,
json_response=True,
token_verifier=token_verifier,
auth=AuthSettings(
issuer_url=AnyHttpUrl(
"https://auth.example.com"
),
resource_server_url=AnyHttpUrl(
"https://mcp.example.com"
),
required_scopes=["mcp:issues.read"],
),
)
@server.tool()
async def get_issue(
issue_id: str,
) -> Issue:
"""Read one issue by its stable ID."""
issue = ISSUES.get(issue_id)
if issue is None:
raise ValueError("issue does not exist")
return issue
return server
这里刻意把 TokenVerifier 作为部署依赖传入,而不是在示例中写一个“只要 Token 非空就通过”的假校验器。真实实现至少要验证:
- Token 签名与允许的算法;
- issuer;
- audience 或 resource;
- 过期时间与生效时间;
- 当前操作要求的 scope;
- Token 是否已被撤销。
如果 Server 背后还要访问 Issue API,应使用 Server 自己取得的下游凭据,或者经过明确的 Token Exchange。不能把 Harness 发来的 MCP Access Token 原样转发给下游 API,这种 token passthrough 会让下游服务无法判断 Token 原本签发给谁。
required_scopes 提供了 Server 入口的最低门槛。更细的授权仍然要在工具实现中完成,例如确认当前主体能否读取指定项目,而不是通过了 mcp:issues.read 就能读取所有租户的数据。
远程 Server 配置仍然不能来自模型
与上一章的 McpServerConfig 一样,远程 endpoint、resource URL、scope 和工具 allowlist 都属于 Harness Policy:
from dataclasses import dataclass
@dataclass(frozen=True)
class RemoteMcpServerConfig:
id: str
endpoint: str
resource_url: str
redirect_uri: str
scopes: tuple[str, ...]
allowed_tools: frozenset[str]
supported_protocol_versions: frozenset[str]
request_timeout_seconds: float = 30.0
stream_read_timeout_seconds: float = 300.0
connect_timeout_seconds: float = 5.0
queue_timeout_seconds: float = 10.0
max_concurrency: int = 8
terminate_on_close: bool = True
endpoint 是真正收发 MCP 消息的地址,例如:
https://mcp.example.com/mcp
resource_url 是 OAuth 中标识目标 Resource Server 的规范地址,例如:
https://mcp.example.com
二者可能相同,也可能不同,不能通过删除 /mcp 猜出 resource。Harness 应保存管理员确认过的值,并检查 Server 发布的 Protected Resource Metadata 与配置一致。
还要在加载配置时做基础 URL 校验:
import ipaddress
from urllib.parse import urlsplit
def validate_remote_url(
value: str,
*,
allow_loopback_http: bool = False,
) -> None:
parsed = urlsplit(value)
if parsed.username or parsed.password:
raise ValueError("URL userinfo is not allowed")
if parsed.query or parsed.fragment:
raise ValueError(
"query and fragment are not allowed"
)
if parsed.hostname is None:
raise ValueError("URL must contain a host")
try:
port = parsed.port
except ValueError as exc:
raise ValueError("URL contains an invalid port") from exc
if parsed.scheme == "https":
return
if not allow_loopback_http or parsed.scheme != "http":
raise ValueError("remote MCP URLs must use HTTPS")
try:
address = ipaddress.ip_address(parsed.hostname)
is_loopback = address.is_loopback
except ValueError:
is_loopback = parsed.hostname == "localhost"
if not is_loopback:
raise ValueError(
"plain HTTP is only allowed on loopback"
)
if port is None:
raise ValueError(
"loopback HTTP must use an explicit port"
)
def validate_remote_config(
config: RemoteMcpServerConfig,
*,
development: bool = False,
) -> None:
validate_remote_url(
config.endpoint,
allow_loopback_http=development,
)
validate_remote_url(
config.resource_url,
allow_loopback_http=development,
)
validate_remote_url(
config.redirect_uri,
allow_loopback_http=development,
)
if config.max_concurrency < 1:
raise ValueError("max_concurrency must be positive")
if not config.allowed_tools:
raise ValueError("allowed_tools cannot be empty")
if not config.scopes or any(
not scope.strip() for scope in config.scopes
):
raise ValueError("scopes cannot be empty")
timeout_values = (
config.request_timeout_seconds,
config.stream_read_timeout_seconds,
config.connect_timeout_seconds,
config.queue_timeout_seconds,
)
if any(value <= 0 for value in timeout_values):
raise ValueError("timeouts must be positive")
if (
config.stream_read_timeout_seconds
< config.request_timeout_seconds
):
raise ValueError(
"stream read timeout cannot be shorter "
"than request timeout"
)
这段校验只负责配置层的明显错误,并不是完整 SSRF 防护。OAuth discovery 会继续访问 Server 返回的 metadata URL、Authorization Server URL 和 token endpoint。攻击者可能把它们指向云 metadata、内网管理服务,或者先解析到公网地址再切换 DNS。
生产环境应按照 MCP 安全最佳实践,通过成熟的 URL/地址校验组件、受控 DNS 和 egress proxy 限制网络目的地,并在每次 redirect 后重新校验。不要用几行字符串判断代替网络层策略。
模型只应该选择已经注册的 server_id:
config = state.remote_mcp_servers[server_id]
它不能提供完整 URL,更不能让工具调用临时修改 scope、redirect URI 或 OAuth issuer。
凭据必须按租户和用户分区
一个共享 MCP Server 不等于一条连接可以服务所有人。至少应使用以下维度划分认证分区:
tenant_id
+ user_id
+ server_id
+ resource_url
= credential partition
同一个用户访问两个 Server,Token 不能混用;同一个 Server 服务两个用户,也不能共享 OAuth storage。
先定义分区与 Secret Store 接口:
import hashlib
from dataclasses import dataclass
from typing import Any, Protocol
@dataclass(frozen=True)
class CredentialPartition:
tenant_id: str
user_id: str
server_id: str
resource_url: str
def storage_prefix(self) -> str:
raw = "\0".join(
(
self.tenant_id,
self.user_id,
self.server_id,
self.resource_url,
)
)
digest = hashlib.sha256(
raw.encode("utf-8")
).hexdigest()
return f"mcp-oauth/{digest}"
class SecretStore(Protocol):
async def read_json(
self,
key: str,
) -> dict[str, Any] | None:
...
async def write_json(
self,
key: str,
value: dict[str, Any],
) -> None:
...
storage_prefix() 的摘要只是避免把用户标识直接放进 key,不是对 Token 的加密。加密、访问控制、密钥轮换和审计由 SecretStore 的真实实现负责。
官方 SDK 的 TokenStorage 可以接到这层 Secret Store:
from mcp.client.auth import TokenStorage
from mcp.shared.auth import (
OAuthClientInformationFull,
OAuthToken,
)
class SecretBackedTokenStorage(TokenStorage):
def __init__(
self,
secret_store: SecretStore,
partition: CredentialPartition,
) -> None:
prefix = partition.storage_prefix()
self.secret_store = secret_store
self.tokens_key = f"{prefix}/tokens"
self.client_key = f"{prefix}/client"
async def get_tokens(self) -> OAuthToken | None:
payload = await self.secret_store.read_json(
self.tokens_key
)
if payload is None:
return None
return OAuthToken.model_validate(payload)
async def set_tokens(
self,
tokens: OAuthToken,
) -> None:
await self.secret_store.write_json(
self.tokens_key,
tokens.model_dump(
mode="json",
exclude_none=True,
),
)
async def get_client_info(
self,
) -> OAuthClientInformationFull | None:
payload = await self.secret_store.read_json(
self.client_key
)
if payload is None:
return None
return OAuthClientInformationFull.model_validate(
payload
)
async def set_client_info(
self,
client_info: OAuthClientInformationFull,
) -> None:
await self.secret_store.write_json(
self.client_key,
client_info.model_dump(
mode="json",
exclude_none=True,
),
)
Token 和动态注册得到的 Client 信息都不能进入:
- system prompt;
messages;- Task checkpoint;
- Artifact 内容;
- 普通应用日志;
- 模型可调用的文件工具。
Task checkpoint 只保存 credential_partition_id 或连接配置 ID,恢复时再由受信任的 Secret Store 取回凭据。
OAuth 交互属于可信 UI
Authorization Code Flow 需要把授权地址交给用户,并接收 redirect callback。Harness 不应该让模型负责这一步,也不应该把授权 URL 塞进 Shell 命令打开。
先定义一个与具体 Web UI 解耦的接口:
from typing import Protocol
class OAuthInteraction(Protocol):
async def present_authorization_url(
self,
partition: CredentialPartition,
authorization_url: str,
) -> None:
...
async def wait_for_callback(
self,
partition: CredentialPartition,
) -> tuple[str, str | None]:
...
实现通常由 Harness 的 Web 控制面完成:
- 为当前登录用户创建一次性 OAuth transaction;
- 将 transaction 绑定 tenant、user、server 和过期时间;
- 在可信 UI 中展示或跳转到授权地址;
- callback 回到固定的 Harness redirect URI;
- 将
code与state交回等待中的 OAuth Client; - transaction 使用一次后立即失效。
同一 credential partition 的交互式授权流程还应该串行执行,并用一次性 transaction ID 关联 UI 与 callback。不能只拿 partition 当 callback key,否则同一用户同时打开两次授权时,后返回的 code 可能唤醒错误的等待者。
然后创建官方 SDK 的 OAuthClientProvider:
from mcp.client.auth import OAuthClientProvider
from mcp.shared.auth import OAuthClientMetadata
from pydantic import AnyUrl
def build_oauth_provider(
config: RemoteMcpServerConfig,
partition: CredentialPartition,
storage: SecretBackedTokenStorage,
interaction: OAuthInteraction,
) -> OAuthClientProvider:
async def redirect_handler(
authorization_url: str,
) -> None:
await interaction.present_authorization_url(
partition,
authorization_url,
)
async def callback_handler() -> tuple[str, str | None]:
return await interaction.wait_for_callback(
partition
)
return OAuthClientProvider(
server_url=config.resource_url,
client_metadata=OAuthClientMetadata(
client_name="Agent Harness",
redirect_uris=[
AnyUrl(config.redirect_uri),
],
grant_types=[
"authorization_code",
"refresh_token",
],
response_types=["code"],
token_endpoint_auth_method="none",
scope=" ".join(config.scopes),
),
storage=storage,
redirect_handler=redirect_handler,
callback_handler=callback_handler,
)
官方 Python SDK 会处理 Protected Resource Metadata 与 Authorization Server metadata discovery、PKCE、state、Token 获取和 refresh。Harness 仍然要限制 discovery 的网络出口,并确保 authorization URL 只能通过可信的 HTTP/HTTPS 导航流程打开,不能传给 sh -c、open 字符串拼接或模型控制的浏览器脚本。
每次授权和 Token 请求都应该携带目标 resource,让签发的 Token 只面向当前 MCP Server。SDK 会按 MCP 授权流程处理它,Server 仍要验证收到的 Token 确实为自身签发。
建立远程 MCP Connection
上一章的 McpConnection 使用 stdio_client() 管理子进程。本章保留同样的外部接口:
connection.config
connection.session
connection.catalog_dirty
因此工具发现和 Binding 逻辑可以直接复用。
import asyncio
import hashlib
from collections.abc import Callable
from contextlib import AsyncExitStack
from datetime import timedelta
import httpx
from mcp import ClientSession
from mcp.client.streamable_http import (
streamable_http_client,
)
from mcp.types import (
Implementation,
ServerNotification,
ToolListChangedNotification,
)
class RemoteMcpConnection:
def __init__(
self,
config: RemoteMcpServerConfig,
oauth_provider: OAuthClientProvider,
) -> None:
self.config = config
self.oauth_provider = oauth_provider
self.catalog_dirty = asyncio.Event()
self.call_slots = asyncio.Semaphore(
config.max_concurrency
)
self._stack: AsyncExitStack | None = None
self._session: ClientSession | None = None
self._get_session_id: (
Callable[[], str | None] | None
) = None
@property
def session(self) -> ClientSession:
if self._session is None:
raise RuntimeError(
"remote MCP connection is not open"
)
return self._session
def session_fingerprint(self) -> str | None:
if self._get_session_id is None:
return None
session_id = self._get_session_id()
if session_id is None:
return None
return hashlib.sha256(
session_id.encode("utf-8")
).hexdigest()[:12]
async def _handle_message(
self,
message: object,
) -> None:
if (
isinstance(message, ServerNotification)
and isinstance(
message.root,
ToolListChangedNotification,
)
):
self.catalog_dirty.set()
async def __aenter__(
self,
) -> "RemoteMcpConnection":
if self._stack is not None:
raise RuntimeError(
"remote MCP connection is already open"
)
stack = AsyncExitStack()
self._stack = stack
try:
timeout = httpx.Timeout(
connect=(
self.config.connect_timeout_seconds
),
read=(
self.config.stream_read_timeout_seconds
),
write=(
self.config.request_timeout_seconds
),
pool=self.config.queue_timeout_seconds,
)
max_connections = (
self.config.max_concurrency + 4
)
http_client = (
await stack.enter_async_context(
httpx.AsyncClient(
auth=self.oauth_provider,
timeout=timeout,
limits=httpx.Limits(
max_connections=max_connections,
max_keepalive_connections=(
max_connections
),
keepalive_expiry=30.0,
),
follow_redirects=False,
headers={
"User-Agent": (
"agent-harness/0.1.0"
),
},
)
)
)
read_stream, write_stream, get_session_id = (
await stack.enter_async_context(
streamable_http_client(
self.config.endpoint,
http_client=http_client,
terminate_on_close=(
self.config
.terminate_on_close
),
)
)
)
session = await stack.enter_async_context(
ClientSession(
read_stream,
write_stream,
read_timeout_seconds=timedelta(
seconds=(
self.config
.request_timeout_seconds
),
),
message_handler=self._handle_message,
client_info=Implementation(
name="agent-harness",
version="0.1.0",
),
)
)
initialization = await session.initialize()
protocol_version = str(
initialization.protocolVersion
)
if (
protocol_version
not in (
self.config
.supported_protocol_versions
)
):
raise McpContractError(
"negotiated protocol version "
"is not allowed"
)
if initialization.capabilities.tools is None:
raise McpContractError(
"server does not advertise tools"
)
self._session = session
self._get_session_id = get_session_id
return self
except BaseException:
self._session = None
self._get_session_id = None
self._stack = None
await stack.aclose()
raise
async def __aexit__(
self,
*exc_info: object,
) -> None:
stack = self._stack
self._session = None
self._get_session_id = None
self._stack = None
if stack is not None:
await stack.aclose()
这里有几个值得单独说明的细节。
follow_redirects=False 是有意设置的安全默认值。若目标 OAuth 部署必须经过 redirect,应在受控 HTTP transport 或 egress proxy 中逐跳校验 scheme、host、DNS 解析结果和目标地址,再为确认过的跳转放行,而不是直接对所有 discovery 请求开启自动跟随。
HTTP timeout 不是一个数字
连接建立、请求写入、普通响应和 SSE 等待的时间尺度不同:
connect应该较短,快速发现 DNS、路由或 TLS 故障;write约束发送请求体的时间;read要允许合法的长时间 SSE stream;pool防止调用无限等待连接池。
即使 HTTP read timeout 很长,单次工具调用仍由 ClientSession.call_tool(..., read_timeout_seconds=...) 设置业务等待上限。两层 timeout 解决的是不同问题。
不手写 MCP Header
streamable_http_client() 会管理:
Accept: application/json, text/event-stream;MCP-Session-Id;- 协商后的
MCP-Protocol-Version; - SSE 的
Last-Event-ID。
Harness 不应再维护另一份 Session ID 状态。get_session_id callback 只用于生成不可逆的观测指纹,不能把原值写进日志。
关闭时尝试终止 Session
terminate_on_close=True 会在存在 Session 时尝试发送 DELETE。Server 可以不支持并返回 405 Method Not Allowed,Client 仍然要正常释放本地资源。
对于“暂停后还要恢复”的长任务,不应该靠不关闭 HTTP Client 来保活。Task checkpoint 保存业务状态;MCP Session 是否能跨进程恢复,要由具体 Server 的协议合同明确保证。
复用上一章的发现与调用代码
上一章的 discover_tools() 和 McpToolAdapter 只依赖:
connection.config.id
connection.config.allowed_tools
connection.config.request_timeout_seconds
connection.session
本地和远程配置都提供这些字段。为了让类型标注也表达这层结构,可以把函数参数从具体 McpConnection 改成 Protocol:
from typing import Protocol
class McpToolConfig(Protocol):
@property
def id(self) -> str:
...
@property
def allowed_tools(self) -> frozenset[str]:
...
@property
def request_timeout_seconds(self) -> float:
...
class McpToolConnection(Protocol):
@property
def config(self) -> McpToolConfig:
...
@property
def session(self) -> ClientSession:
...
然后只调整类型注解:
async def discover_tools(
connection: McpToolConnection,
) -> list[McpToolBinding]:
# 函数体与上一章保持一致
...
McpToolAdapter.__init__() 的 connection 参数也同样改成 McpToolConnection,函数体不变。这样 McpConnection 与 RemoteMcpConnection 都按结构满足同一接口,类型检查不会把远程连接误判成不兼容对象。
运行逻辑没有分叉:
- 完成
initialize; - 分页调用
tools/list; - 检查 allowlist 与 schema;
- 冻结当前模型轮次的 Catalog snapshot;
- 生成本地工具别名;
- 将
tools/call结果写回原来的 Provider tool call ID。
这也说明 Streamable HTTP 应该位于连接层,而不是渗透进 Agent Loop。模型不需要知道这次工具调用最终走的是 JSON response 还是 SSE stream。
为共享 Server 增加背压
本地 Task Server 通常只服务一个 Agent Loop,远程 Server 却可能被很多 Task 同时调用。HTTP 连接池只能限制 socket 数量,不能代替业务并发限制。
在上一章的 Adapter 外增加一个轻量 wrapper:
from typing import Any
class RemoteMcpToolAdapter(McpToolAdapter):
def __init__(
self,
connection: RemoteMcpConnection,
bindings: list[McpToolBinding],
) -> None:
super().__init__(connection, bindings)
self.connection = connection
async def call(
self,
local_name: str,
arguments: dict[str, Any],
) -> str:
try:
async with asyncio.timeout(
self.connection.config
.queue_timeout_seconds
):
await self.connection.call_slots.acquire()
except TimeoutError:
return (
"Error: remote MCP server is busy; "
"retry after other calls finish"
)
try:
return await super().call(
local_name,
arguments,
)
finally:
self.connection.call_slots.release()
这层限制应该按认证分区维护,而不是全局只有一把锁:
(tenant A, user 1, server X) -> connection + semaphore
(tenant A, user 2, server X) -> connection + semaphore
(tenant B, user 1, server X) -> connection + semaphore
不同分区可以配置额外的总量配额,避免某个租户占满整个 Server。一个连接内部是否允许多条并发 tools/call,还要以目标 Server 的实际能力测试为准;如果它依赖严格顺序,就把 max_concurrency 设置为 1。
批量 tool calls 也不能无上限 asyncio.gather()。Provider 一次返回十几个调用时,仍然要经过:
- 当前模型轮次的调用数量上限;
- 每个 Server 的 semaphore;
- 每个 tenant 的 rate limit;
- HTTP pool;
- 单次调用 timeout。
分清流恢复、会话重建与工具重试
远程 MCP 中“重连”至少可能指三件不同的事:
| 动作 | 保留什么 | 是否重新发送 tools/call |
|---|---|---|
| SSE 流恢复 | 同一 Session、同一逻辑请求 | 不重新发送 |
| MCP 会话重建 | 新 Session、新初始化与目录快照 | 不自动重放在途调用 |
| 工具重试 | 新的一次业务调用 | 会重新发送 |
SSE 流恢复
如果 Server 为 SSE event 提供 id,Client 可以在 stream 中断后使用 Last-Event-ID 继续接收遗漏消息。官方 SDK 会在当前 transport 内处理这类 resumption。
这不是重新调用工具。它只是继续读取同一次逻辑请求的后续事件,因此不会因为恢复 stream 而生成第二个 Provider tool result。
MCP 会话重建
如果 Server 返回 404 表示 Session 已终止,旧的 MCP-Session-Id 不能继续使用。Harness 应:
- 停止把新调用发给旧连接;
- 将所有在途写操作标记为
unknown; - 完整关闭旧
streamable_http_clientcontext; - 创建新的 HTTP transport 和
ClientSession; - 重新执行
initialize; - 重新发现、校验并冻结 Tool Catalog;
- 在下一个安全的模型轮次原子替换 Registry。
不能只对旧 ClientSession 再调用一次 initialize()。Transport 内部仍然持有已经失效的 Session ID,必须从连接 context 开始重建。
工具重试
工具重试是一个业务决定,与 transport 重连没有必然关系。
例如 create_issue 已经在 Server 执行成功,但返回前 Session 失效。新建连接后再次调用,会创建第二条 Issue。除非调用携带由 Harness 生成并由 Server 强制执行的业务幂等键,否则不能自动重放写操作。
恢复策略可以沿用第 12 章:
| 情况 | Harness 行为 |
|---|---|
| 本地 schema 校验失败 | 写回模型,允许修正参数 |
isError=true | 写回模型,连接继续使用 |
| SSE 短暂中断且可恢复 | 由 transport 按 event ID 恢复,不重放调用 |
| Session 终止 | 重建连接与目录,在途副作用标记为 unknown |
401 Unauthorized | refresh 或重新授权,不把 Token 暴露给模型 |
403 insufficient_scope | 经用户确认后有限度 step-up,限制授权重试次数 |
| timeout / 网络中断 | 不自动重试写操作,先使用查询工具核对副作用 |
输出违反 outputSchema | 记录 contract failure,不把结果交给模型 |
若任务必须支持取消,Harness 应保存 MCP request 与 Job Record 的关联,并发送协议取消 notification。即便 Server 确认收到取消,也要由工具合同说明外部副作用是否能够回滚。
一条完整的远程启动路径
把认证、连接、发现和 Registry 串起来:
from contextlib import asynccontextmanager
from dataclasses import replace
from typing import AsyncIterator
@asynccontextmanager
async def remote_task_state(
state: HarnessState,
config: RemoteMcpServerConfig,
partition: CredentialPartition,
secret_store: SecretStore,
interaction: OAuthInteraction,
) -> AsyncIterator[HarnessState]:
validate_remote_config(config)
if config.id != partition.server_id:
raise ValueError(
"credential partition does not match server"
)
if config.resource_url != partition.resource_url:
raise ValueError(
"credential partition does not match resource"
)
storage = SecretBackedTokenStorage(
secret_store,
partition,
)
oauth_provider = build_oauth_provider(
config,
partition,
storage,
interaction,
)
async with RemoteMcpConnection(
config,
oauth_provider,
) as connection:
bindings = await discover_tools(connection)
adapter = RemoteMcpToolAdapter(
connection,
bindings,
)
task_registry = state.tool_registry.fork()
install_mcp_tools(task_registry, adapter)
yield replace(
state,
tool_registry=task_registry,
)
调用方在 context 内运行 Task:
async def run_remote_agent(
state: HarnessState,
config: RemoteMcpServerConfig,
partition: CredentialPartition,
secret_store: SecretStore,
interaction: OAuthInteraction,
user_message: str,
) -> str:
async with remote_task_state(
state,
config,
partition,
secret_store,
interaction,
) as task_state:
return await run_agent_loop(
task_state,
user_message,
)
这份结构保持了上一章建立的生命周期约束:
- Registry 中的 remote handler 不会逃出连接 context;
- 工具在认证、初始化和目录校验完成前不会暴露给模型;
- 连接关闭后,Session Registry 一起丢弃;
- 下次重建会生成新的目录快照;
- OAuth Token 由 Secret Store 持有,不进入
HarnessState的可序列化部分。
实际系统可以用 Connection Manager 在多个 Task 间复用同一认证分区的健康连接,但必须使用引用计数或 lease 管理生命周期。不要让一个 Task 退出时关闭另一个 Task 正在使用的连接,也不要跨用户复用连接来节省一次 TLS 握手。
Origin、CORS 与 SSRF 是三件事
这三个名词经常一起出现在远程 MCP 部署中,但保护的方向不同。
Origin 校验保护 Server
Streamable HTTP Server 必须校验请求中的 Origin。如果请求携带了无效 Origin,应返回 403。这可以降低恶意网页利用本机或内网 MCP Server 的 DNS rebinding 风险。
后端 httpx Client 通常不会发送 Origin,不需要凭空伪造一个。Server 应正确处理“没有 Origin 的非浏览器 Client”,并对“存在 Origin 的浏览器请求”执行 allowlist。
CORS 控制浏览器读取响应
如果浏览器直接连接 MCP Server,需要配置精确的 allowed origins、methods 和 headers;若要读取 Session header,还要暴露:
MCP-Session-Id
使用 Cookie 或 Authorization header 时,不要使用 Access-Control-Allow-Origin: *。CORS 也不是认证,Server 仍要验证每个请求的 Token。
SSRF 保护 Harness
OAuth discovery 让 Harness 主动访问远端 metadata 和 token endpoint,因此 Harness 自己也可能成为 SSRF 发起者。防护应覆盖:
- MCP endpoint;
- Protected Resource Metadata;
- Authorization Server metadata;
- registration、authorization、token 与 redirect URL;
- DNS 解析后的最终地址;
- 每一次 HTTP redirect。
最可靠的边界是配置 allowlist 加受控 egress,而不是相信 Server 返回的 URL。
观测连接,但不要记录凭据
远程故障需要足够信息才能定位。一次工具调用可以记录:
trace_id
task_id
server_id
tenant_fingerprint
session_fingerprint
catalog_snapshot_id
provider_tool_call_id
remote_tool_name
connection_generation
queue_wait_ms
duration_ms
result_class
不应记录:
Authorization header
access token
refresh token
authorization code
raw MCP-Session-Id
OAuth client secret
完整的敏感工具参数与结果
session_fingerprint() 只用于判断两条日志是否来自同一 Session,不能用于恢复。还要检查 HTTP Client 和 MCP SDK 的日志级别与 redaction filter,避免底层库在 debug 或 info 日志中输出 header、URL query 或原始 Session ID。
建议至少建立以下指标:
| 指标 | 用途 |
|---|---|
| connection open / close / rebuild | 判断连接是否频繁抖动 |
| OAuth refresh / reauthorization | 发现 Token 生命周期和授权配置问题 |
| queue wait | 判断并发上限是否过低或 Server 过载 |
| call latency / timeout | 区分慢工具与网络问题 |
| session terminated | 判断 Server 状态是否稳定 |
| catalog refresh / contract failure | 发现 Server 升级导致的 schema 变化 |
| unknown side effect | 追踪需要人工或补偿流程核对的操作 |
对 unknown 写操作要生成持久化 Recovery Record,而不是只打一行 error 日志。Harness 重启后仍然需要知道哪些外部副作用尚未确认。
常见坑
第一,把旧 HTTP+SSE 当成 Streamable HTTP。 新协议使用一个支持 POST 与 GET 的 MCP endpoint。
第二,认为所有返回都必须是 SSE。 短请求可以直接返回 JSON;是否流式由 Server 和请求决定。
第三,把 Session ID 当成登录凭据。 每个请求仍要验证 Access Token,并确认 Session 与主体绑定。
第四,只在 initialize 时认证一次。 HTTP 请求彼此独立,后续 POST、GET 和 DELETE 都要鉴权。
第五,让模型提供 MCP URL。 endpoint、resource、issuer、redirect URI 和 scope 必须来自 Harness 配置。
第六,所有用户共享一份 TokenStorage。 OAuth 凭据至少按 tenant、user、server 和 resource 分区。
第七,把 Token 写进 Task checkpoint。 checkpoint 保存凭据引用,真正的 Token 只进入 Secret Store。
第八,使用 Session ID 恢复业务任务。 Session 只保存协议状态,业务进度仍由 Task、Job 和 Artifact 持久化。
第九,Session 失效后复用旧 transport。 关闭旧 context,重新创建 transport、ClientSession、初始化与目录快照。
第十,把 SSE 断开视为取消。 断线不会自动停止 Server 执行,需要协议取消与业务补偿。
第十一,断线后自动重放写工具。 超时和 Session 终止都不能证明工具没有执行。
第十二,只用 HTTP connection pool 控制并发。 还需要 Server、tenant 和认证分区级的 semaphore 与 quota。
第十三,盲目开启 HTTP redirects。 OAuth discovery 与 redirect 都可能成为 SSRF 路径,应逐跳校验并限制 egress。
第十四,认为 CORS 等于认证。 CORS 只约束浏览器读取,Token 验证和工具授权仍不可省略。
第十五,把 Client Token 转发给下游 API。 MCP Server 必须验证 Token,并使用独立下游凭据或明确的 Token Exchange。
第十六,在日志里打印原始 Session 和 Authorization header。 使用指纹与脱敏字段完成关联。
第十七,连接重建后沿用旧 Tool Catalog。 新 Session 重新发现工具,并在模型轮次边界替换快照。
第十八,无状态工具也默认开启 stateful session。 先使用 stateless JSON,确有流式和跨请求状态需求时再增加 Session。
小结
本文把上一章的 MCP Client 从本地子进程延伸到了远程共享服务:
- Streamable HTTP 使用统一 endpoint,每条 Client 消息通过 POST 发送;
- Server 可以返回 JSON 或 SSE,也可以通过 GET 建立主动消息流;
- SSE 断开不等于工具取消,流恢复也不等于重新调用工具;
- 简单工具优先使用 stateless HTTP 与 JSON response;
- 受保护的 MCP Server 是 OAuth Resource Server,认证服务独立存在;
- Access Token、MCP Session 和 Provider tool call ID 属于不同层次;
- endpoint、resource、scope 和 allowlist 全部来自 Harness Policy;
- OAuth Token 按 tenant、user、server 和 resource 隔离,并保存到 Secret Store;
- 远程 Connection 复用上一章的初始化、工具发现、Binding 与 Schema 校验;
- HTTP pool、调用 semaphore、tenant quota 和 timeout 共同形成背压;
- SSE resumption 保留原请求,Session 失效则重建连接与目录;
- 在途写操作遇到断线或超时后进入
unknown,不能盲目重放; - Origin、CORS 和 SSRF 保护不同方向,不能互相替代;
- Trace 记录连接代次、目录快照和不可逆指纹,不记录 Token 与原始 Session。
至此,Agent 可以用同一套 Tool Registry 连接 Task 内的本地 MCP Server,也可以连接团队共享的远程 MCP Server。传输位置变了,但权限、审批、Tool Catalog、Task、Job 与恢复边界没有被绕开。
下一章将对前面的所有内容做一次系统总结:从最初的 Agent Loop 出发,重新串联 Tool Registry、权限与 Hook、TodoWrite、Subagent、Skill、上下文压缩、Memory、Prompt Assembler、错误恢复、Task、后台任务、Scheduler、Agent Team、Worktree 与 MCP,并把这些模块放回一张完整架构图和一条端到端执行链路中。我们也会进一步区分哪些能力构成 Agent Harness 的最小核心,哪些属于面向长任务、多 Agent 协作和生产环境的工程增强。