深入 MCP Python SDK 的 Subscriptions:用 subscriptions/listen 流式订阅服务器变更
【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk
服务端的目录(catalog)从来不是一成不变的:工具会在运行时动态出现,资源 URI 背后的内容会持续变化。订阅(Subscriptions)正是客户端获知这些变化的机制——客户端发送一次subscriptions/listen请求,而该请求的响应本身就是一条流(stream):它保持打开,持续携带客户端所请求的那些变更通知。本文以i18n/hi/pages/handlers/subscriptions.md(英文原文见 docs/handlers/subscriptions.md)为主线,结合 python-sdk 仓库的服务端实现(src/mcp/server/subscriptions.py)与配套测试(tests/docs_src/test_subscriptions.py),完整讲解服务端的订阅机制:如何发布变更、过滤契约、访问控制、客户端配合、跨进程扩展以及低层(low-level)组合方式。读完你将掌握一套可直接运行的、从单进程到多副本的订阅发布方案。
订阅机制概览:一次请求,一条响应流
在 2026-07-28 协议(subscriptions/listen,SEP-2575)下,订阅是"请求即流"的形态:
- 客户端发送
subscriptions/listen请求,携带一个filter,声明它关心哪几类变更; - 服务端把响应当作一条常开的流,先发送确认帧(acknowledgment),随后凡是命中 filter 的变更通知都沿这条流下发;
- 流一直保持打开,直到服务端主动优雅关闭,或客户端断开。
SDK 的MCPServer已经为你内置了subscriptions/listen的完整服务能力。从源码看,src/mcp/server/mcpserver/server.py 中,MCPServer在构造 low-levelServer时自动挂载了on_subscriptions_listen=ListenHandler(self._subscriptions);wire 层面的职责——首帧确认、逐流过滤、每个帧上打订阅 id 标记——全部由 SDK 承担,业务代码只需要"发布变更"这一个动作。
从工具里发布变更:一行代码
服务端要做的只有一件事:把变更 publish 出去。完整示例见 docs_src/subscriptions/tutorial001.py,核心代码如下:
from mcp.server.mcpserver import Context, MCPServer mcp = MCPServer("Sprint Board") BOARDS = { "sprint": {"design": False, "build": False, "ship": False}, "backlog": {"tidy docs": False}, } @mcp.resource("board://{name}") def board(name: str) -> str: tasks = BOARDS[name] return "\n".join(f"[{'x' if done else ' '}] {task}" for task, done in tasks.items()) @mcp.tool() async def complete_task(board: str, task: str, ctx: Context) -> str: BOARDS[board][task] = True await ctx.notify_resource_updated(f"board://{board}") return f"{task}: done" def sprint_report() -> str: done = sum(done for tasks in BOARDS.values() for done in tasks.values()) return f"{done} task(s) done" @mcp.tool() async def enable_reports(ctx: Context) -> str: mcp.add_tool(sprint_report) await ctx.notify_tools_changed() return "reporting is live"逐行拆解这里的发布语义:
await ctx.notify_resource_updated("board://sprint"):只抵达每一个对该 URI 订阅了的打开流,其他流一概收不到。await ctx.notify_tools_changed():抵达所有请求了"工具列表变更"的流。收到该事件的客户端会重新调用tools/list,此刻就能看到新注册的sprint_report。- 两个姊妹方法:
notify_prompts_changed()与notify_resources_changed(),分别对应 prompt 列表与资源列表的变更。 - 没有订阅者就不干活:对空闲服务器 publish 是纯 no-op。因此你永远不必检查"有没有人在听",只需要陈述"什么东西变了"。
从源码看,这四个notify_*方法定义于 src/mcp/server/mcpserver/context.py:它们本质上是把对应的 typed 事件(ToolsListChanged、PromptsListChanged、ResourcesListChanged、ResourceUpdated(uri=...))投递到Context持有的SubscriptionBus上。事件词汇表定义在 src/mcp/shared/subscriptions.py,服务端与客户端共享同一套 typed 事件,绝不直接传递 JSON-RPC。
wire 上的真实形态
文档用!!! check展示了一条 filter 中写明了board://sprint的流,在complete_task执行后的真实 wire 数据:
{"method": "notifications/subscriptions/acknowledged", "params": {"notifications": {"resourceSubscriptions": ["board://sprint"]}, "_meta": {"io.modelcontextprotocol/subscriptionId": "listen-1"}}} {"method": "notifications/resources/updated", "params": {"uri": "board://sprint", "_meta": {"io.modelcontextprotocol/subscriptionId": "listen-1"}}}注意更新通知里没有携带 board 内容——它只带 URI。每个帧都在_meta下携带 listen 请求的 JSON-RPC id,这个 id 就是subscription id。id 由客户端铸造:Python 的Client使用"listen-1"这类字符串,其他客户端可能用整数。SDK 侧把_meta的键常量定义为SUBSCRIPTION_ID_META_KEY = "io.modelcontextprotocol/subscriptionId"(见 src/mcp/shared/subscriptions.py)。此外,确认帧是流的第一帧——ListenHandler在订阅总线之后、开始排空事件缓冲之前发送 ack(src/mcp/server/subscriptions.py),这样即使 ack 写入被挂起期间有事件发布,也会被缓冲而不会丢失,且 ack 依然稳稳占据首帧位置。
过滤即契约:只送达被请求的内容
filter 是一份契约。一个同时请求了工具列表变更和某一个资源 URI 的流,只会收到这两类东西,其他一概不发——哪怕你发布了 prompt 变更,那条流也保持沉默。
两个必须明确的匹配细节:
- URI 按精确字符串匹配:
MCPServer把资源 URI 当作逐字字符串来比对,所以订阅了board://sprint的流,听不到任何关于board://sprint/tasks/1的消息。规范允许服务端上报"订阅 URI 的子资源发生变化",但MCPServer从不这么做;不过客户端被设计为对这种行为有所预期(读取event.uri而非假定具体哪个资源动了,见 docs/client/subscriptions.md)。 - 在 src/mcp/shared/subscriptions.py 中,
event_matches是服务端投递与客户端接收共用的准入谓词:ResourceUpdated只有event.uri in uris(honored 资源 URI 集合)时才命中,三类ListChanged事件则要求 filter 中对应标志为True。_honored_subset(src/mcp/server/subscriptions.py)还会把非True的标志与空 URI 列表从确认帧中剔除,而不是回显 falsy 值。
这条流不是什么
- 不是回放日志(replay log):断掉的流就永远没了,无人连接期间发布的事件不会被排队。客户端重新 listen、重新 fetch 即可。
ListenHandler的投递是 fire-and-forget 且无回放的(src/mcp/server/subscriptions.py)。 - 不是 2025 年的老路径:调用过
resources/subscribe的旧客户端,由ctx.session.send_resource_updated(uri)服务(服务端会话 API)。notify_*系列方法只抵达subscriptions/listen流。两者在语义与协议代际上是分开的。
决定谁可以观看:用 middleware 给订阅上闸门
默认情况下,每个被请求的 kind 与 URI 都会被接受:任何调用者都能 watch 你发布的任何 URI。这里没有任何机制咨询你的 read handler——因为根本没人"读"。一个会被你的files://{name}handler 拒之门外的调用者,依然可以对着files://payroll.csv开一条流,从而得知"它变了、以及什么时候变的"。他永远拿不到内容,也探测不出存在性——因为未知 URI 同样会被 honor,只是永远不会触发。这个风险虽窄但真实存在,所以在多租户服务器上按用户发布 URI 之前,务必给它加闸门。
这道闸门是一个middleware。它在 SDK 确认subscriptions/listen请求之前先看到请求,一旦调用者请求了其无权读取的内容就拒绝。完整示例见 docs_src/subscriptions/tutorial006.py:
from mcp_types import INVALID_REQUEST, SubscriptionsListenRequestParams from mcp.server.auth.middleware.auth_context import get_access_token from mcp.server.context import CallNext, HandlerResult, ServerRequestContext from mcp.server.mcpserver import MCPServer from mcp.shared.exceptions import MCPError # Who may see each file. Replace this table with a database or your RBAC system. ACCESS = { "files://report.pdf": {"alice", "bob"}, "files://payroll.csv": {"carol"}, } def can_access(user: str | None, uri: str) -> bool: return user is not None and user in ACCESS.get(uri, set()) async def gate_subscriptions(ctx: ServerRequestContext, call_next: CallNext) -> HandlerResult: if ctx.method == "subscriptions/listen": params = SubscriptionsListenRequestParams.model_validate(ctx.params or {}, by_name=False) token = get_access_token() user = token.subject if token else None if not all(can_access(user, uri) for uri in params.notifications.resource_subscriptions or ()): raise MCPError(INVALID_REQUEST, "not permitted to watch the requested resources") return await call_next(ctx) mcp = MCPServer("Reports", middleware=[gate_subscriptions]) @mcp.resource("files://{name}") def file(name: str) -> str: uri = f"files://{name}" token = get_access_token() if not can_access(token.subject if token else None, uri): raise MCPError(INVALID_REQUEST, f"Unknown resource: {uri}") return f"contents of {name}"实现要点:
ctx.params是原始请求,所以 middleware 自己用SubscriptionsListenRequestParams.model_validate(...)校验,并读出客户端请求的 filter。- 拒绝 = 在
call_next(ctx)之前抛出MCPError:客户端收到这个错误、得不到任何流,而连接照常存活。错误消息要保持统一、不点名任何 URI,这样"被拒绝"这件事永远不能反向确认哪些 URI 是受保护的。 - 同一个
can_access(user, uri)回答两个问题:资源 handler 在resources/read上问它,middleware 在subscriptions/listen上问它。把示例中的表换成数据库或你自己的 RBAC 系统,两条路径依然同步。 - 决定在整个流生命周期内生效:事件送达时不会逐条复查。如果调用者的访问权限可能在流中途失效(例如 token 过期),请在失效时主动断开该调用者的连接。
middleware 的完整契约、它还包裹哪些内容、为何被标记为 provisional,见 docs/advanced/middleware.md。配套测试test_the_middleware_refuses_a_listen_the_caller_could_not_read(tests/docs_src/test_subscriptions.py)验证了:alice 可以 watchfiles://report.pdf,而对files://payroll.csv的 listen 会在任何确认之前被 in-band 拒绝,且resources/read同样被拒。
客户端那一端:async with client.listen(...)
与上述流相对的另一端,是盯着 board 的客户端,完整代码见 docs_src/subscriptions/tutorial003.py:
from mcp import Client from mcp.client.subscriptions import ResourceUpdated, ToolsListChanged from mcp.types import TextResourceContents BOARD = "board://sprint" async def read_board(client: Client, uri: str = BOARD) -> str: [contents] = (await client.read_resource(uri)).contents assert isinstance(contents, TextResourceContents) return contents.text async def follow_board(client: Client) -> None: async with client.listen(tools_list_changed=True, resource_subscriptions=[BOARD]) as sub: async for event in sub: match event: case ResourceUpdated(uri=uri): print(await read_board(client, uri)) case ToolsListChanged(): tools = await client.list_tools() print("tools:", [tool.name for tool in tools.tools]) case _: pass # kinds the filter did not ask for never arrive async def main() -> None: async with Client("http://localhost:8000/mcp") as client: await follow_board(client)关键语义:
- 进入
client.listen(...)的瞬间请求就已发出,并等待你的确认帧,因此 block 一开始,流就是活的;之后发布的每一个变更都不会漏掉。 - 每个 typed 事件都是"重新 fetch"的信号,绝不是 payload:
ResourceUpdated只带uri,客户端据此重新read_resource;ToolsListChanged驱动重新list_tools。这正是测试test_follow_board_prints_the_refetched_board_and_the_new_tool_list(tests/docs_src/test_subscriptions.py)验证的行为——事件驱动重取,board 重印、新工具名重列。 - 客户端细节(主流程旁并行监听、流结束处理、重新 listen)在Clients目录下的独立页面 docs/client/subscriptions.md 中有完整讲述,包括
sub.honored、sub.subscription_id、SubscriptionLost异常与 1024 条未消费事件上限等。本页只交代与服务器侧对接的"全貌"。
扩展到一个进程之外:实现 SubscriptionBus
发布动作从你的 handler 到打开的流之间,经由SubscriptionBus传递。默认实现是进程内的:一个进程、其中的每条流。在没有负载均衡副本之前,这就是正确答案;一旦你跑多副本,客户端的流会钉在某个 replica 上,而发生在另一个 replica 上的 publish 必须能到达它——这时就需要换 bus。
这个接缝由你来实现:在你的 pub/sub 后端之上实现两个方法。文档给出基于 Redis 的完整示例:
from collections.abc import Callable from redis.asyncio import Redis from mcp.server.mcpserver import MCPServer from mcp.server.subscriptions import ServerEvent # SubscriptionBus is a Protocol: no base class class RedisSubscriptionBus: def __init__(self, redis: Redis) -> None: self._redis = redis self._listeners: dict[object, Callable[[ServerEvent], None]] = {} async def publish(self, event: ServerEvent) -> None: await self._redis.publish("mcp-events", encode(event)) # to every replica def subscribe(self, listener: Callable[[ServerEvent], None]) -> Callable[[], None]: token = object() self._listeners[token] = listener def unsubscribe() -> None: self._listeners.pop(token, None) return unsubscribe mcp = MCPServer("Sprint Board", subscriptions=RedisSubscriptionBus(redis))对照源码契约理解这份实现:
SubscriptionBus是一个Protocol(src/mcp/server/subscriptions.py),没有基类,只有两个方法:异步的publish(event)与同步的subscribe(listener) -> unsubscribe。publish设计为异步,是为了让后端实现可以做网络 I/O。encode是你自己的,每个 replica 上的 reader task 也由你负责:它解码到达的消息并调用每个已注册的 listener。listeners 是同步的、不允许抛异常、运行在服务器的事件循环上。内置的InMemorySubscriptionBus也是同样的语义:单个 listener 抛异常只会被记录并跳过,不会饿死其他 listener(src/mcp/server/subscriptions.py)。- bus 上跑的是 typed
ServerEvent,四个小型 dataclass,绝不出现 JSON-RPC。stamping(打订阅 id 标记)、filtering(逐流过滤)和流生命周期都留在 SDK 内,因此任何 bus 实现都不可能破坏协议——它只能把事件在进程间搬运。事件匹配与 wire 通知的相互转换由event_matches/event_to_notification(src/mcp/shared/subscriptions.py)完成。 - 为了在请求之外发布,请自己构造 bus 以持有引用。什么都不传时,
MCPServer会在内部构建一个InMemorySubscriptionBus且不对外暴露(src/mcp/server/mcpserver/server.py)。自己持有时,就能从 lifespan 任务、webhook 等任何地方发布:
from mcp.server.subscriptions import InMemorySubscriptionBus, ToolsListChanged bus = InMemorySubscriptionBus() mcp = MCPServer("Sprint Board", subscriptions=bus) async def tools_reloaded() -> None: await bus.publish(ToolsListChanged()) # from a lifespan task, a webhook, anywhere低层组合:三条线把同样的部件装起来
在低层Server上没有任何预接线,同样的部件三条线即可组装完成。完整示例见 docs_src/subscriptions/tutorial002.py:
from typing import Any import mcp.types as types from mcp.server.context import ServerRequestContext from mcp.server.lowlevel import Server from mcp.server.subscriptions import InMemorySubscriptionBus, ListenHandler, ResourceUpdated bus = InMemorySubscriptionBus() listen_handler = ListenHandler(bus) BOARD = {"design": False, "build": False} COMPLETE_TASK_SCHEMA: dict[str, Any] = { "type": "object", "properties": {"task": {"type": "string"}}, "required": ["task"], } async def read_resource( ctx: ServerRequestContext[Any], params: types.ReadResourceRequestParams ) -> types.ReadResourceResult: board = "\n".join(f"[{'x' if done else ' '}] {task}" for task, done in BOARD.items()) return types.ReadResourceResult(contents=[types.TextResourceContents(uri=params.uri, text=board)]) async def list_tools( ctx: ServerRequestContext[Any], params: types.PaginatedRequestParams | None ) -> types.ListToolsResult: return types.ListToolsResult( tools=[types.Tool(name="complete_task", description="Mark a task done.", input_schema=COMPLETE_TASK_SCHEMA)] ) async def call_tool(ctx: ServerRequestContext[Any], params: types.CallToolRequestParams) -> types.CallToolResult: args = params.arguments or {} BOARD[args["task"]] = True await bus.publish(ResourceUpdated(uri="board://sprint")) return types.CallToolResult(content=[types.TextContent(type="text", text="done")]) server = Server( "sprint-board", on_read_resource=read_resource, on_list_tools=list_tools, on_call_tool=call_tool, on_subscriptions_listen=listen_handler, )三个要点:
- bus 归你所有,所以你可以直接对它 publish:
await bus.publish(ResourceUpdated(uri=...))。把它放在 handler 够得着的地方——这里是模块作用域,更大的应用里放在 lifespan 中。 ListenHandler(bus)正是MCPServer注册的那个 handler,on_subscriptions_listen=只是一个普通 handler 槽位。想换语义,就在那个槽里放你自己的 callable——但那样规范义务就转移给你了:先确认、每个帧都打上订阅 id、filter 之外的东西一律不送。ListenHandler的两个内建护栏也值得注意:max_subscriptions=1024限制并发流数(超限在确认前以INTERNAL_ERROR拒绝),max_buffered_events=1024限制每条流的积压上限(积压见顶即结束该流,客户端重新 listen 即可,因为无回放,结束流不会丢失积压本就在丢的东西),见 src/mcp/server/subscriptions.py。ListenHandler.close()优雅结束每条打开的流:每条流都会把 listen 请求的 result 作为最后一帧收到——这正是规范中"服务端故意结束订阅"的表达方式。该方法在流 flush 完成之前就返回,所以拆除 transport 之前要给它们留一点时间。不调用它的话,流会在客户端断开时才结束。
低层组合与高层MCPServer用的是同一套机制:测试test_lowlevel_composition_serves_the_same_stream(tests/docs_src/test_subscriptions.py)证明两条路径产出的流行为一致——确认帧先行、精确 URI 过滤、tools_list_changed之外的 kind 不会送达。
总结
- 客户端用一次
subscriptions/listen请求加入,响应就是流;服务端提供它开箱即用。 - 你用
ctx.notify_*发布,stamping、filtering 与生命周期的工作由 SDK 完成。 - 事件是信号(cue),不是 payload:两端收到信号后各自重新 fetch。
- 客户端那一端就是
async with client.listen(...),其完整故事见Clients下的 docs/client/subscriptions.md。 - 低层
Server上你自己组装同样的部件:一个 bus、ListenHandler(bus)、on_subscriptions_listen槽位。 - 扩容 = 实现
SubscriptionBus,只需两个方法,然后以MCPServer(subscriptions=...)传入。
运行承载这一切的服务器——单副本也好、二十副本也罢——见 docs/run/deploy.md。
【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址: https://gitcode.com/gh_mirrors/pythonsd/python-sdk
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考