:从 notify_* 发布到 SubscriptionBus 跨进程扩展)
人工智能MCP 服务MCP Clients【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址https://gitcode.com/gh_mirrors/pythonsd/python-sdk点击查看免费下载导读本文聚焦 Model Context Protocol Python SDK当前仓库gh_mirrors/pythonsd/python-sdk中基于 2026-07-28 协议时代SEP-2575的subscriptions/listen订阅机制。服务端的目录并非一成不变——工具会在运行时出现资源 URI 背后的内容会变化而订阅正是客户端获知这些变化的方式。读完本文你将掌握如何在工具处理函数中用一行ctx.notify_*发布变更、订阅在线上wire的真实帧格式与过滤器契约、如何用中间件做订阅授权、如何从低层Server手工组合出同一套订阅能力以及如何通过实现SubscriptionBus协议将订阅扩展到多副本部署。原文主体见 docs/handlers/subscriptions.md含多语言译本 i18n/fr/pages/handlers/subscriptions.md配套代码见 docs_src/subscriptions/。订阅的本质一次请求其响应就是流MCPServer的目录是动态的工具在运行时被添加、资源 URI 背后的内容会变化。在 2026-07-28 协议时代客户端不再依赖常驻的 GET 流如旧版 SSE而是发送一条subscriptions/listen请求——这条请求的响应本身就是流它保持打开持续承载客户端所请求的那些变更通知。从源码看src/mcp/shared/subscriptions.py定义了服务端与客户端共享的类型化事件词汇表四个小 dataclassToolsListChanged—— 工具列表发生变化PromptsListChanged—— 提示prompt列表发生变化ResourcesListChanged—— 资源列表发生变化ResourceUpdated(uri...)—— 指定uri的资源内容变化可能需要重新读取。它们共同构成ServerEvent联合类型。模块注释明确写道每个事件都是电平触发level trigger语义是“这个变了在意就重新拉取”因此两端都能通过去重来约束缓冲区。这四个事件在服务端被src/mcp/server/subscriptions.py重新导出并统一收敛为ServerEvent别名。从工具中发布变更你的工作只有一行服务端发布变更的代码量被压缩到一行。以官方教程 docs_src/subscriptions/tutorial001.py 中的 “Sprint Board” 服务器为例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(fboard://{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 的开放流除此之外无人收到。它的实现位于 src/mcp/server/mcpserver/context.py内部即await self._bus.publish(ResourceUpdated(uristr(uri)))——URI 会先被str()归一化再进入总线。await ctx.notify_tools_changed()触达每一个请求了工具列表变更的流。客户端收到该事件后会再次调用tools/list此时便能看到新注册的sprint_report。其实现同样是对ToolsListChanged()的一次publishcontext.py。姊妹方法还有notify_prompts_changed()与notify_resources_changed()分别发布PromptsListChanged与ResourcesListChangedcontext.py。没有订阅者就没有工作向空闲服务器发布是一个 no-op因此你永远不需要检查是否有人在监听只需声明“什么变了”即可。MCPServer会替你自动注册并服务subscriptions/listen这一方法低层Server则需手动传入on_subscriptions_listen见后文。所有线上义务——以确认帧acknowledgment作为流的第一帧、按流过滤、在每一帧上打订阅 ID 标记——都是 SDK 的职责而不是你的。线上的样子确认帧、订阅 ID 与“信号而非载荷”教程中给出了complete_task执行后、过滤器包含board://sprint的流在线上会看到的两帧 JSON{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 的内容。事件只是“信号”两端都会据此重新拉取数据。每一帧都在_meta下携带io.modelcontextprotocol/subscriptionId键其值就是那条listen请求的 JSON-RPC id也就是订阅 ID。该键名在源码中被定义为常量SUBSCRIPTION_ID_META_KEY io.modelcontextprotocol/subscriptionId见 src/mcp/shared/subscriptions.py。这个 ID 是客户端铸造的Python 的Client使用listen-1这样的字符串其他客户端可能使用整数因此_meta的值类型是通用的。服务端的实现细节与文档叙述完全对应ListenHandler.__call__先从总线注册自己再发送确认帧——这样当确认写操作被挂起时发布的事件会被缓冲而不是丢失确认仍是第一帧因为只有该 handler 任务写这条流并且它在确认发送返回之后才开始消费缓冲区见 src/mcp/server/subscriptions.py。过滤器是契约只交付被请求的内容过滤器是一份契约。一个流请求了工具列表变更加一个资源 URI就只会收到这两种事件其他一律不交付——你发布一个 prompt 变更该流保持沉默。MCPServer对资源 URI 做精确字符串匹配命名了board://sprint的流听不到任何关于board://sprint/tasks/1的消息。这一点在Context.notify_resource_updated的 docstring 中写得很清楚context.py过滤谓词event_matches(honored, uris, event)也确实是event.uri in uris的集合成员判断src/mcp/shared/subscriptions.py。规范允许服务器报告某个已订阅 URI 的子资源发生变化MCPServer从不这么做但客户端被设计为对这种情况有所准备。同时这个流不是两样东西它不是重放日志。流一旦中断就彻底丢失没有人连接期间发布的事件也不会被排队。客户端需要重新 listen 并重新拉取数据。它不是 2025 时代的路径。调用过resources/subscribe的旧客户端由ctx.session.send_resource_updated(uri)服务notify_*方法只触达subscriptions/listen的流。谁可以观察用中间件在订阅入口做授权默认情况下每个被请求的种类和 URI 都会被兑现任何调用者都可以观察你发布的任何 URI。注意——没有任何代码会查询你的读取 handler因为根本没有人“读取”。一个会被你的files://{name}handler 拒绝的调用者依然可以为files://payroll.csv打开流并得知它变了、何时变的。它永远学不到内容也无法探测哪些 URI 存在未知 URI 同样会被“兑现”只是永远不触发。这个漏洞很窄但真实存在在多租户服务器上发布按用户区分的 URI 之前务必加上控制。控制手段是中间件。它会在 SDK 发送确认之前看到subscriptions/listen请求并在调用者请求了任何无权读取的内容时拒绝。完整示例见 docs_src/subscriptions/tutorial006.pyfrom 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_nameFalse) 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 ffiles://{name} token get_access_token() if not can_access(token.subject if token else None, uri): raise MCPError(INVALID_REQUEST, fUnknown resource: {uri}) return fcontents of {name}要点拆解ctx.params是原始请求因此中间件要自己用SubscriptionsListenRequestParams.model_validate(...)把它校验成类型化结构再读出客户端请求的过滤器params.notifications.resource_subscriptions。拒绝方式是在call_next(ctx)之前抛出MCPError客户端拿到这个错误、得不到任何流而连接继续存活。错误消息要保持统一、不要点名任何 URI这样一次拒绝永远不会“确认”哪些 URI 受保护。一个can_access(user, uri)函数同时回答两个问题资源 handler 在resources/read上问它中间件在subscriptions/listen上问它。把示例里的 ACCESS 表换成数据库或你的 RBAC 系统两处自动保持同步。决定覆盖流的整个生命周期没有逐事件复查。如果调用者的访问权可能在流中途失效比如令牌过期请在失效时主动结束该调用者的连接。关于中间件的完整契约它还包裹了哪些东西、为什么被标记为 provisional见 docs/advanced/middleware.md。客户端视角async with client.listen(...)下面是这条流另一端、跟随同一个 board 的客户端docs_src/subscriptions/tutorial003.pyfrom 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_changedTrue, resource_subscriptions[BOARD]) as sub: async for event in sub: match event: case ResourceUpdated(uriuri): 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(...)会发出请求并等待服务器的确认因此当代码块开始时流已经激活每个类型化事件都是一次重新拉取的信号绝不是载荷。客户端listen的入口实现位于 src/mcp/client/subscriptions.py其参数tools_list_changed、prompts_list_changed、resources_list_changed、resource_subscriptions就是订阅过滤器。客户端一侧还有更多内容——在主流程旁并行观察、流的结束与重听、SubscriptionLost异常丢失的流无法重放观察者需要anyio.sleep退避后重新 listen 并重新拉取见 docs_src/subscriptions/tutorial005.py以及sub.honored服务器确认的过滤器与sub.subscription_id等句柄属性——这些都在 docs/client/subscriptions.md 的Clients章节下详细展开本文不再赘述。跨进程扩展实现SubscriptionBus协议发布事件从你的 handler 到达开放流靠的是SubscriptionBus。默认总线在内存中一个进程、它包含的所有流。在单进程下这是正确答案但一旦你在负载均衡器后面运行多个副本客户端流会被钉在某个副本上而另一个副本上的发布必须能到达它。这个接缝由你来实现只需在 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, subscriptionsRedisSubscriptionBus(redis))encode由你实现每个副本上解码入站消息并调用每个已注册 listener 的读取任务也是你的活。从源码看SubscriptionBus被定义为Protocol而非基类src/mcp/server/subscriptions.py因此无需继承协议约束只有两个方法async def publish(self, event: ServerEvent)—— 异步以便后端实现可以执行网络 I/Odef subscribe(self, listener) - Callable[[], None]—— 同步的本地注册返回幂等的退订 callable。其余约束体现在文档与实现注释中listener 是同步的、不得抛异常、运行在服务器的事件循环上内存总线会在 fan-out 边界捕获并记录单个 listener 的异常避免一个坏 listener 饿死其他 listener 或拖垮发布 handler见 subscriptions.py。总线运送的是类型化的ServerEvent值——那四个小 dataclass绝不是 JSON-RPC。打标记stamping、过滤和流生命周期都留在 SDK 内所以一个总线实现不可能破坏协议它只能把事件在进程之间搬运。在请求之外发布比如 lifespan 任务或 webhook 里时需要自己构造总线以持有引用——MCPServer在你不传参时会在内部构建一个默认总线且不对外暴露from mcp.server.subscriptions import InMemorySubscriptionBus, ToolsListChanged bus InMemorySubscriptionBus() mcp MCPServer(Sprint Board, subscriptionsbus) async def tools_reloaded() - None: await bus.publish(ToolsListChanged()) # from a lifespan task, a webhook, anywhere低层组合没有预接线三行拼装在低层Server上没有任何预接线同样的部件用三行就能拼起来完整代码见 docs_src/subscriptions/tutorial002.pyfrom mcp.server.subscriptions import InMemorySubscriptionBus, ListenHandler, ResourceUpdated bus InMemorySubscriptionBus() listen_handler ListenHandler(bus) # ... 资源读取、工具列表、工具调用等 handler ... server Server( sprint-board, on_read_resourceread_resource, on_list_toolslist_tools, on_call_toolcall_tool, on_subscriptions_listenlisten_handler, )总线归你所有因此你可以直接向它发布await bus.publish(ResourceUpdated(uri...))。把它放在 handler 能触达的地方这里放在模块作用域更大的应用里放在 lifespan 中。ListenHandler(bus)与MCPServer内部注册的是同一个 handler而on_subscriptions_listen只是一个普通的 handler 槽位。如果你想实现不同的语义可以把自定义 callable 放进这个槽位——但那样规格义务就回到你身上先确认、给每一帧打上订阅 ID、绝不交付过滤器之外的内容。ListenHandler.close()会优雅地结束每一个开放流每条流都会收到 listen 请求的结果作为最后一帧——这正是规范中“服务器有意结束订阅”的表达方式。它返回时流可能还没冲刷完毕所以要给它们一点时间再拆卸传输层否则流会在客户端断开连接时结束。从实现看ListenHandler还带两个可调参数max_subscriptions默认 1024限制并发流数超出时在确认前以INTERNAL_ERROR拒绝和max_buffered_events默认 1024限制每条流的事件积压积压达到上限的流会被结束——客户端重听并重新拉取即可因为没有重放结束流并不会丢失积压之外的东西见 src/mcp/server/subscriptions.py。小结客户端用一条subscriptions/listen请求完成订阅响应本身就是流MCPServer开箱即用地服务该方法。你用ctx.notify_*发布notify_resource_updated、notify_tools_changed、notify_prompts_changed、notify_resources_changed打标记、过滤、生命周期全部由 SDK 负责。事件是信号而非载荷两端收到后都重新拉取数据。客户端一侧就是async with client.listen(...)完整故事在 docs/client/subscriptions.md 的Clients章节。在低层Server上自己拼装同样的部件一个总线、ListenHandler(bus)、on_subscriptions_listen槽位。水平扩展意味着实现SubscriptionBus两个方法并通过MCPServer(subscriptions...)传入。运行这套服务、无论后面是一个副本还是二十个副本都是 docs/run/deploy.mdDéployer et passer à léchelle的主题服务端订阅的授权、过滤与发布细节则始终以本文与 docs/handlers/subscriptions.md 为准。赞分享人工智能MCP 服务MCP Clients【免费下载链接】python-sdkThe official Python SDK for Model Context Protocol servers and clients项目地址https://gitcode.com/gh_mirrors/pythonsd/python-sdk点击查看免费下载相关推荐UVR Ultimate Vocal Remover 如何快速分离人声与伴奏UVR Ultimate Vocal Remover 如何快速分离人声与伴奏 Ultimate Vocal Remover UVR 是一款基于深度学习的图形界面人工智能MCP 服务MCP Clients深入 MCP Python SDK 的 Subscriptions用 subscriptions/listen 流式订阅服务器变更深入 MCP Python SDK 的 Subscriptions用 subscriptions/listen 流式订阅服务器变更 服务端的目录catalo人工智能MCP 服务MCP ClientsMCP Python SDK 订阅机制完全指南从 notify_* 发布、listen 流到多副本扩展MCP Python SDK 订阅机制完全指南从 notify_ 发布、listen 流到多副本扩展 服务端的目录catalog并非一成不变工具会在运行人工智能MCP 服务MCP Clients创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考