Essay

Agent Harness 工程:实现 MCP Server 与 Client

By XiaoLeiJun

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。

Agent Harness 中的 MCP 角色与边界

模型只生成普通工具调用;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 提供的模板。

它们拥有不同的协议方法,不应该为了省代码全部伪装成工具。

还有三组名字尤其容易混淆:

  1. MCP 的实验性 Tasks 不等于第 13 章的持久化 Task Graph;
  2. MCP Roots 用来告诉 Server 有哪些文件根,不是文件系统沙箱;
  3. 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。

MCP 初始化、工具发现与调用时序

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 制造的竞态。

更稳定的策略是:

  1. 模型请求开始前冻结 Tool Catalog snapshot;
  2. 当前回合的所有调用都使用同一快照;
  3. 回合结束后检查 catalog_dirty
  4. 重新分页发现、校验和过滤;
  5. 原子替换下一轮 Registry;
  6. 在 Trace 与 checkpoint 中记录 snapshot ID。

模型看见的工具集合,必须在一次决策周期里保持稳定。

MCP Server 的生命周期跟随 Task Worktree

本地 Workspace Server 最清楚的粒度是:

一个 Task
-> 一个 Worktree
-> 一个 MCP Server 进程
-> 一个 Client Session
-> 一份 Task Tool Registry

Task Session 启动时:

  1. 根据 Worktree 构造 McpServerConfig
  2. 在任务沙箱中启动 Server;
  3. 完成初始化;
  4. 发现并安装允许的工具;
  5. 把 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,让传输跨过网络,但不让信任边界跟着变得模糊。