Agent Harness 工程:实现 MCP Server 与 Client
Agent Harness 工程:实现 MCP Server 与 Client
到上一章为止,Agent 已经拥有 Tool Registry、权限、Task、Job 和独立 Worktree。
可所有工具仍然住在 Harness 进程里:
Agent Loop -> Tool Registry -> Python function
文件工具是 Python 函数,Issue 系统要写专用 HTTP Client,数据库再写一套连接代码。每多接一种能力,发现、参数校验、超时和错误处理就多长出一条分支。
更麻烦的是,工具升级或崩溃也会直接影响 Agent 主进程。
我们希望把能力移到进程边界之外,同时保留一个共同入口:
Agent Loop
-> Tool Registry
-> MCP Client
<=> MCP Server
-> real capability
MCP,也就是 Model Context Protocol,为 Host 与能力提供者规定了一套共同语言。
但协议不会自动替我们完成权限、沙箱和恢复。这一章真正要做的,是把 MCP 接到 Harness 的正确位置。
先看清谁在和谁说话
MCP 的名字里有 Model,但模型通常不会直接连接 Server。
模型只生成普通工具调用;Harness 管理权限与上下文,MCP Client 负责协议通信,Server 执行真实能力。
四个角色分别是:
Model
根据可见 schema 提出 tool_call
Agent Harness / Host
管理 Agent Loop、Task、权限、Worktree 和结果
MCP Client
与一个 Server 协商协议、发现能力并转发请求
MCP Server
暴露 Tools、Resources 或 Prompts
一个 Host 可以同时维护多个 Client,每个 Client 对应一个 Server。Client 是 Harness 内部的协议组件,不是另一个 Agent。
本章只接入 Tools。
- Tools 是可执行动作;
- Resources 是由 URI 标识的上下文数据;
- Prompts 是 Server 提供的模板。
它们拥有不同的协议方法,不应该为了省代码全部伪装成工具。
还有三组名字尤其容易混淆:
- MCP 的实验性 Tasks 不等于第 13 章的持久化 Task Graph;
- MCP Roots 用来告诉 Server 有哪些文件根,不是文件系统沙箱;
- Tool annotation 是提示,不是 Harness 可以直接相信的权限事实。
MCP 统一了能力交换,安全责任仍然属于 Host。
固定规范,也固定 SDK 范围
本文按 MCP 2025-11-25 规范编写,并固定 Python SDK 的主版本范围:
pip install "mcp>=1.27,<2" "jsonschema>=4.20,<5"
协议版本和 Python 包版本是两套编号。协议可能保持兼容,SDK API 却会随主版本调整,所以示例需要上界,项目还应把精确解析版本写进 lockfile。
不要自己手写 JSON-RPC 帧、请求 ID、初始化通知和进程退出逻辑。官方 SDK 已经负责这些协议细节,我们只在外层管理信任边界和运行时状态。
第一个 Server 先从 stdio 开始
MCP 2025-11-25 定义了两种标准传输:
stdio
Host 在本地启动受控子进程
Streamable HTTP
连接远程或共享 Server
本章先用 stdio,因为它与 Task Worktree 很自然地配合:
创建 Task Worktree
-> 以 Worktree 为 cwd 启动 Server
-> Server 跟随 Task Session
-> Session 结束时关闭进程
stdio 有一条看似细小、实际非常严格的规则:
stdin / stdout 只传 MCP 消息
普通日志只能写 stderr
Server 里随手 print("started"),对人类是日志,对 Client 却是一条无法解析的协议消息。
写一个只读当前 Worktree 的 Server
先创建 workspace_server.py:
import os
from pathlib import Path
from typing import Annotated
from mcp.server.fastmcp import FastMCP
from pydantic import BaseModel, Field
MAX_FILE_BYTES = 2_000_000
WORKSPACE_ROOT = Path(
os.environ["AGENT_WORKSPACE_ROOT"]
).resolve(strict=True)
mcp = FastMCP(
"workspace-tools",
instructions=(
"Tools operate inside the assigned workspace."
),
)
class FilePreview(BaseModel):
path: str
content: str
truncated: bool
def resolve_workspace_path(
relative_path: str,
) -> Path:
candidate = Path(relative_path)
if candidate.is_absolute():
raise ValueError(
"absolute paths are not allowed"
)
resolved = (
WORKSPACE_ROOT / candidate
).resolve(strict=True)
if (
resolved != WORKSPACE_ROOT
and WORKSPACE_ROOT not in resolved.parents
):
raise ValueError(
"path escapes task workspace"
)
return resolved
@mcp.tool()
def read_file(
path: str,
line_limit: Annotated[
int,
Field(ge=1, le=1000),
] = 200,
) -> FilePreview:
"""Read one UTF-8 file in the assigned workspace."""
target = resolve_workspace_path(path)
if not target.is_file():
raise ValueError(
"path is not a regular file"
)
if target.stat().st_size > MAX_FILE_BYTES:
raise ValueError("file is too large")
lines: list[str] = []
truncated = False
with target.open(
"r",
encoding="utf-8",
) as handle:
for line_number, line in enumerate(handle):
if line_number >= line_limit:
truncated = True
break
lines.append(line.rstrip("\r\n"))
return FilePreview(
path=str(
target.relative_to(WORKSPACE_ROOT)
),
content="\n".join(lines),
truncated=truncated,
)
if __name__ == "__main__":
mcp.run(transport="stdio")
FastMCP 根据函数签名生成 inputSchema。返回 Pydantic Model 时,它还会生成 outputSchema,并提供结构化结果。
工具抛出 ValueError 时,Server 返回一次 isError=true 的工具结果,而不是让整条连接崩溃。模型可以根据错误修正路径再试。
路径校验能阻止绝对路径、.. 和已有 symlink 的明显越界,但它仍然不是完整沙箱。Server 进程必须运行在上一章的 Task Sandbox 中,只挂载当前 Worktree。
Server 怎么启动,不能由模型决定
模型不能提供可执行文件、工作目录和环境变量。它只调用 Harness 已经批准的能力。
from dataclasses import dataclass
from pathlib import Path
@dataclass(frozen=True)
class McpServerConfig:
id: str
command: str
args: tuple[str, ...]
cwd: Path
env: dict[str, str]
allowed_tools: frozenset[str]
supported_protocol_versions: frozenset[str]
request_timeout_seconds: float = 30.0
当前 Task 的配置由 Harness 从 Worktree 派生:
import sys
def build_workspace_server_config(
workspace: TaskWorkspace,
server_script: Path,
) -> McpServerConfig:
root = workspace.root.resolve(strict=True)
runtime_home = (
workspace.temp_dir.resolve(strict=True)
/ "home"
)
runtime_home.mkdir(parents=True, exist_ok=True)
return McpServerConfig(
id=f"workspace-{workspace.task_id}",
command=sys.executable,
args=(
str(server_script.resolve(strict=True)),
),
cwd=root,
env={
"AGENT_WORKSPACE_ROOT": str(root),
"HOME": str(runtime_home),
"PYTHONUNBUFFERED": "1",
},
allowed_tools=frozenset({"read_file"}),
supported_protocol_versions=frozenset(
{"2025-11-25"}
),
request_timeout_seconds=30,
)
这里有四层约束:
- command 和脚本来自部署配置;
- cwd 与环境中的 Workspace Root 绑定当前 Task;
allowed_tools是拒绝优先的本地 allowlist;- HOME 指向 Task Runtime Root,不暴露用户目录配置。
环境过滤仍然不能替代沙箱。网络、凭据文件和系统调用要由 OS 边界继续限制。
第一条交互必须是初始化
MCP 连接不能一启动就调用 tools/list。
根据生命周期规范,Client 与 Server 先完成:
Client -> initialize
Server -> protocolVersion + capabilities + serverInfo
Client -> notifications/initialized
官方 ClientSession.initialize() 会完成初始化与通知。Harness 再检查协商版本和必需能力。
import asyncio
from contextlib import AsyncExitStack
from datetime import timedelta
from mcp import ClientSession, StdioServerParameters
from mcp.client.stdio import stdio_client
from mcp.types import (
Implementation,
ServerNotification,
ToolListChangedNotification,
)
class McpContractError(RuntimeError):
pass
class McpConnection:
def __init__(
self,
config: McpServerConfig,
) -> None:
self.config = config
self.catalog_dirty = asyncio.Event()
self._stack: AsyncExitStack | None = None
self._session: ClientSession | None = None
@property
def session(self) -> ClientSession:
if self._session is None:
raise RuntimeError(
"MCP connection is not open"
)
return self._session
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,
) -> "McpConnection":
if self._stack is not None:
raise RuntimeError(
"MCP connection is already open"
)
stack = AsyncExitStack()
self._stack = stack
try:
parameters = StdioServerParameters(
command=self.config.command,
args=list(self.config.args),
cwd=self.config.cwd,
env=dict(self.config.env),
)
read_stream, write_stream = (
await stack.enter_async_context(
stdio_client(parameters)
)
)
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(
"protocol version is not allowed"
)
if (
initialization.capabilities.tools
is None
):
raise McpContractError(
"server does not advertise tools"
)
self._session = session
return self
except BaseException:
self._session = None
self._stack = None
await stack.aclose()
raise
async def __aexit__(
self,
*exc_info: object,
) -> None:
stack = self._stack
self._session = None
self._stack = None
if stack is not None:
await stack.aclose()
AsyncExitStack 让子进程、stdio 流和 ClientSession 按相反顺序关闭。不要再额外维护一套与 SDK 竞争的 subprocess 清理逻辑。
连接还应该由同一个 Task Session 打开和关闭。底层异步 task group 与取消域需要稳定所有者,不能在一个任务中 __aenter__(),再由另一个任务随意退出。
初始化完成,才轮到发现工具
initialize 只说明 Server 支持 Tools,不会直接返回目录。Client 还要调用 tools/list。
工具发现要同时处理:
- cursor 分页;
- 重复远端名称;
- 本地 allowlist;
- JSON Schema 合法性;
- 多 Server 名称冲突;
- 目录大小上限。
每个远端工具先变成一条 Binding:
import hashlib
import re
from copy import deepcopy
from dataclasses import dataclass
from typing import Any
from jsonschema import Draft202012Validator
from jsonschema.protocols import Validator
from jsonschema.validators import validator_for
def build_validator(
schema: dict[str, Any],
) -> Validator:
validator_class = validator_for(
schema,
default=Draft202012Validator,
)
validator_class.check_schema(schema)
return validator_class(schema)
def make_local_tool_name(
server_id: str,
remote_name: str,
) -> str:
source = f"mcp_{server_id}_{remote_name}"
slug = re.sub(
r"[^A-Za-z0-9_-]+",
"_",
source,
).strip("_")
digest = hashlib.sha256(
f"{server_id}\0{remote_name}".encode()
).hexdigest()[:8]
return f"{slug[:51]}_{digest}"
@dataclass(frozen=True)
class McpToolBinding:
local_name: str
server_id: str
remote_name: str
description: str
input_schema: dict[str, Any]
output_schema: dict[str, Any] | None
input_validator: Validator
output_validator: Validator | None
def provider_schema(self) -> dict[str, Any]:
return {
"type": "function",
"function": {
"name": self.local_name,
"description": self.description,
"parameters": deepcopy(
self.input_schema
),
},
}
模型看到的是稳定本地别名,Harness 保存 local-to-remote 映射。不同 Server 即使都有 read_file,也不会互相覆盖。
这里还要注意一层兼容性:MCP 工具可以使用完整 JSON Schema,但模型 Provider 往往只接受其中一个子集。Harness 应在暴露工具前检查兼容性;无法无损转换的工具要明确拒绝,不能悄悄删掉约束,否则模型看到的合同会比 Server 真正执行的合同更宽松。
完整分页发现可以保持克制:
from mcp.types import PaginatedRequestParams
MAX_TOOL_PAGES = 100
MAX_DISCOVERED_TOOLS = 1000
async def discover_tools(
connection: McpConnection,
) -> list[McpToolBinding]:
bindings: list[McpToolBinding] = []
remote_names: set[str] = set()
seen_cursors: set[str] = set()
cursor: str | None = None
for _ in range(MAX_TOOL_PAGES):
params = (
PaginatedRequestParams(cursor=cursor)
if cursor is not None
else None
)
page = await connection.session.list_tools(
params=params
)
for tool in page.tools:
if tool.name in remote_names:
raise McpContractError(
f"duplicate tool: {tool.name}"
)
remote_names.add(tool.name)
if len(remote_names) > MAX_DISCOVERED_TOOLS:
raise McpContractError(
"tool catalog is too large"
)
if (
tool.name
not in connection.config.allowed_tools
):
continue
input_schema = deepcopy(tool.inputSchema)
if input_schema.get("type") != "object":
raise McpContractError(
f"{tool.name} needs object input"
)
output_schema = (
deepcopy(tool.outputSchema)
if tool.outputSchema is not None
else None
)
bindings.append(
McpToolBinding(
local_name=make_local_tool_name(
connection.config.id,
tool.name,
),
server_id=connection.config.id,
remote_name=tool.name,
description=(
tool.description
or f"MCP tool {tool.name}"
)[:2000],
input_schema=input_schema,
output_schema=output_schema,
input_validator=build_validator(
input_schema
),
output_validator=(
build_validator(output_schema)
if output_schema is not None
else None
),
)
)
cursor = page.nextCursor
if cursor is None:
break
if cursor in seen_cursors:
raise McpContractError(
"tool cursor repeated"
)
seen_cursors.add(cursor)
else:
raise McpContractError(
"tool pagination limit exceeded"
)
missing = (
connection.config.allowed_tools
- remote_names
)
if missing:
raise McpContractError(
f"configured tools missing: {sorted(missing)}"
)
connection.catalog_dirty.clear()
return bindings
发现回答“Server 声称有什么”,allowlist 回答“Harness 允许什么”,模型本轮可见集合还可以更小。
这三层不要混在一起:
Server catalog
-> deployment allowlist
-> current Task tool snapshot
一次工具调用跨过两套协议
模型返回的是 Provider 的:
tool_call_id + local_name + arguments
MCP SDK 则在 Client 与 Server 之间维护自己的 JSON-RPC request ID。
Provider tool_call_id 负责模型消息闭环,MCP request ID 只属于连接内部。
二者不能混为一谈。模型不会生成 MCP JSON-RPC,Server 也不需要知道 Provider 的消息格式。
适配器负责输入校验、调用和输出归一化:
import json
from datetime import timedelta
from jsonschema.exceptions import ValidationError
from mcp.types import TextContent
def read_text_blocks(
blocks: list[object],
) -> str:
texts: list[str] = []
for block in blocks:
if not isinstance(block, TextContent):
raise McpContractError(
"non-text content needs an artifact"
)
texts.append(block.text)
return "\n".join(texts)
class McpToolAdapter:
def __init__(
self,
connection: McpConnection,
bindings: list[McpToolBinding],
) -> None:
self.connection = connection
self.bindings = {
item.local_name: item
for item in bindings
}
if len(self.bindings) != len(bindings):
raise McpContractError(
"local tool name collision"
)
async def call(
self,
local_name: str,
arguments: dict[str, Any],
) -> str:
binding = self.bindings.get(local_name)
if binding is None:
return f"Error: unknown tool {local_name}"
try:
binding.input_validator.validate(
arguments
)
except ValidationError as exc:
return (
"Error: invalid arguments - "
f"{exc.message}"
)
result = (
await self.connection.session.call_tool(
binding.remote_name,
arguments,
read_timeout_seconds=timedelta(
seconds=(
self.connection.config
.request_timeout_seconds
)
),
)
)
text = read_text_blocks(
list(result.content)
)
if result.isError:
detail = text
if (
not detail
and result.structuredContent is not None
):
detail = json.dumps(
result.structuredContent,
ensure_ascii=False,
)
return (
f"Error: {binding.remote_name} failed"
+ (f"\n{detail}" if detail else "")
)
if binding.output_validator is not None:
if result.structuredContent is None:
raise McpContractError(
"structured output is missing"
)
try:
binding.output_validator.validate(
result.structuredContent
)
except ValidationError as exc:
raise McpContractError(
"structured output is invalid"
) from exc
if result.structuredContent is not None:
return json.dumps(
result.structuredContent,
ensure_ascii=False,
separators=(",", ":"),
)
return text
这里刻意区分三类结果:
参数不符合 inputSchema
普通工具错误,模型可以修正
isError=true
连接正常,远端工具执行失败
连接异常或输出违反 contract
进入 Harness Recovery,不伪装成业务错误
图片、音频和嵌入资源不能静默丢掉,也不应该把大段 base64 塞进 messages。更完整的 Adapter 会把它们写入 Artifact Store,再返回 ID、类型和受控读取方式。
Tool Registry 变成异步,Agent Loop 不必重写
MCP 是异步 I/O,因此统一 Tool Registry 也要支持 await:
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
from functools import partial
AsyncToolHandler = Callable[
[dict[str, Any]],
Awaitable[str],
]
@dataclass
class RegisteredTool:
schema: dict[str, Any]
handler: AsyncToolHandler
class ToolRegistry:
def __init__(self) -> None:
self._tools: dict[str, RegisteredTool] = {}
def register_async(
self,
name: str,
schema: dict[str, Any],
handler: AsyncToolHandler,
) -> None:
if name in self._tools:
raise ValueError(
f"duplicate tool: {name}"
)
self._tools[name] = RegisteredTool(
schema=deepcopy(schema),
handler=handler,
)
def fork(self) -> "ToolRegistry":
child = ToolRegistry()
child._tools = {
name: RegisteredTool(
schema=deepcopy(tool.schema),
handler=tool.handler,
)
for name, tool in self._tools.items()
}
return child
async def call(
self,
name: str,
arguments: dict[str, Any],
) -> str:
tool = self._tools.get(name)
if tool is None:
return f"Error: unknown tool {name}"
return await tool.handler(arguments)
def install_mcp_tools(
registry: ToolRegistry,
adapter: McpToolAdapter,
) -> None:
for binding in adapter.bindings.values():
registry.register_async(
name=binding.local_name,
schema=binding.provider_schema(),
handler=partial(
adapter.call,
binding.local_name,
),
)
fork() 为当前 Task 建立独立的工具目录快照,避免连接关闭后,MCP handler 仍留在全局 Registry 中。本地同步工具可以通过 asyncio.to_thread() 包一层;原生异步工具直接 await。不要在事件循环里同步等待 MCP coroutine。
Hook 的顺序保持不变:
解析参数
-> before_tool Hooks
-> 权限与审批
-> Tool Registry
-> MCP tools/call
-> 输出校验
-> after_tool Hooks
-> Provider tool message
权限必须发生在 MCP 请求发出之前。否则远端副作用已经发生,再拒绝把结果写回模型也没有意义。
工具目录变化,要等到安全点再换
Server 可以发送:
notifications/tools/list_changed
前面的 message handler 只设置 catalog_dirty,没有立即修改 Registry。
假设模型刚看到工具 A 的 schema,Server 通知目录变化,Harness 立刻删除 A。模型随后按自己刚看到的目录调用 A,就会遭遇 Harness 制造的竞态。
更稳定的策略是:
- 模型请求开始前冻结 Tool Catalog snapshot;
- 当前回合的所有调用都使用同一快照;
- 回合结束后检查
catalog_dirty; - 重新分页发现、校验和过滤;
- 原子替换下一轮 Registry;
- 在 Trace 与 checkpoint 中记录 snapshot ID。
模型看见的工具集合,必须在一次决策周期里保持稳定。
MCP Server 的生命周期跟随 Task Worktree
本地 Workspace Server 最清楚的粒度是:
一个 Task
-> 一个 Worktree
-> 一个 MCP Server 进程
-> 一个 Client Session
-> 一份 Task Tool Registry
Task Session 启动时:
- 根据 Worktree 构造
McpServerConfig; - 在任务沙箱中启动 Server;
- 完成初始化;
- 发现并安装允许的工具;
- 把 Registry snapshot 交给 Agent Loop。
结束时反向释放:
停止新模型回合
-> 等待或取消在途工具调用
-> 停止 Job 对连接的引用
-> 关闭 MCP Session 与子进程
-> 保存 Trace、checkpoint 和 Artifact
-> 最后回收 Worktree
如果先删除 Worktree,仍在运行的 Server 会突然失去 cwd 和文件。
模型不需要在 read_file 参数里传 workspace_id。工作区由 Task Session 注入,所以它即使知道另一个 Task ID,也不能把当前 Server 切到对方目录。
超时之后,最危险的是自作聪明地重试
Client 等待超时,只能证明结果没有及时回来,不能证明工具没有执行。
例如远端 create_issue 已经成功,响应却在网络中丢失。Harness 如果自动重试,就会创建第二条 Issue。
恢复策略应该区分:
本地参数校验失败
让模型修正
isError=true
把工具错误交给模型,连接继续使用
输出违反 outputSchema
记录 Server contract failure
stdio Server 退出
在安全点重连并重新发现目录
超时或连接中断
副作用状态记为 unknown,先检查真实结果
只读且由本地 Policy 确认为幂等的工具可以有限重试。写操作需要业务幂等键、Tool Execution Record 或查询型补偿工具。
MCP request ID 不是业务幂等键。连接重建后换一个 JSON-RPC ID,并不会让 Server 认出“这还是刚才那次创建”。
小结
这一章把 Tool Registry 延伸到了进程边界之外:
- Harness 是 MCP Host,模型不直接连接 Server;
- stdio Server 跟随 Task Worktree 与 Session;
- Server 配置、环境和工具 allowlist 来自本地 Policy;
- Client 先初始化,再分页发现工具;
- Binding 保存本地别名、远端名称和 Schema;
- Provider tool call ID 与 MCP request ID 各自属于一层协议;
- inputSchema 在请求前校验,outputSchema 在成功后复核;
isError与连接故障走不同恢复路径;- Hook、审批、Artifact 和 Worktree 仍由 Harness 控制;
- Tool Catalog 按模型回合冻结;
- 连接关闭后,才回收对应工作区。
现在,工具不再必须和 Agent Loop 挤在同一个进程里。
不过,stdio Server 仍然跟随当前 Task。Issue、日历和知识库这类团队共享服务,通常已经运行在网络另一端,还要面对认证、Session、断线和多租户隔离。
下一章把同一套 Client 扩展到 远程 MCP,让传输跨过网络,但不让信任边界跟着变得模糊。