Suscripciones
Traducción automática
Esta página se tradujo automáticamente a partir de la documentación en inglés, y la página en inglés es la versión de referencia. Si algo no se lee bien, Traducciones explica cómo avisarnos.
El catálogo de un servidor no es fijo. Las herramientas aparecen en tiempo de ejecución, y el contenido detrás del URI de un recurso cambia.
Las suscripciones son la forma en que un cliente se entera. El cliente envía una solicitud subscriptions/listen, y la respuesta a esa solicitud es el flujo: queda abierto y transporta las notificaciones de cambio que el cliente pidió.
Publícalo desde la herramienta
Tu parte es una sola línea: publicar el cambio.
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")llega a cada flujo abierto que se suscribió a ese URI. A nadie más.await ctx.notify_tools_changed()llega a cada flujo que pidió los cambios en la lista de herramientas. Un cliente que lo recibe vuelve a llamar atools/list, y ahora vesprint_report.- Los métodos hermanos son
notify_prompts_changed()ynotify_resources_changed(). - Sin suscriptores, sin trabajo. Publicar en un servidor inactivo no hace nada, así que nunca compruebas si alguien está escuchando. Declaras qué cambió.
MCPServer atiende subscriptions/listen por ti. Las obligaciones del canal (el acuse de recibo como primera trama, el filtrado por flujo, el id de suscripción en cada trama) son trabajo del SDK.
Check
En el canal, un flujo cuyo filtro nombró board://sprint se ve así después de que se ejecuta complete_task:
{"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"}}}
Fíjate en lo que la actualización no lleva: el tablero. Cada trama lleva el id JSON-RPC de la solicitud listen bajo _meta, y ese id es el id de suscripción. Lo genera el cliente: el Client de Python usa cadenas como "listen-1"; otros clientes pueden usar enteros.
Solo lo que se pidió
El filtro es un contrato. Un flujo que solicitó los cambios en la lista de herramientas y un URI de recurso recibe esos dos tipos y nada más. Publica un cambio de prompt y ese flujo se queda en silencio.
MCPServer compara los URI de recurso como cadenas exactas, así que un flujo que nombró board://sprint no se entera de nada sobre board://sprint/tasks/1. La especificación permite que un servidor informe de un cambio en un subrecurso de un URI suscrito; MCPServer nunca lo hace, pero los clientes están construidos para esperarlo.
Dos cosas que el flujo no es:
- No es un registro de repetición. Un flujo caído se pierde, y los eventos publicados mientras nadie estaba conectado no se encolan. Los clientes vuelven a escuchar y vuelven a consultar.
- No es la vía de 2025. A los clientes que llamaron a
resources/subscribelos atiendectx.session.send_resource_updated(uri). Los métodosnotify_*llegan solo a los flujos desubscriptions/listen.
Decidir quién puede observar
Por defecto se acepta cada tipo y URI solicitado: cualquier llamador puede observar cualquier URI que publiques. Nada consulta tu handler de lectura, porque nadie está leyendo: un llamador al que tu handler files://{name} rechazaría puede igualmente abrir un flujo sobre files://payroll.csv y enterarse de que cambió, y cuándo. Nunca conoce el contenido, y no puede sondear qué existe, porque un URI desconocido también se acepta y simplemente nunca se dispara. Acotado pero real, así que contrólalo antes de publicar URI por usuario desde un servidor multiinquilino.
El control es un middleware. Ve la solicitud subscriptions/listen antes de que el SDK la acuse y la rechaza cuando el llamador pide algo que no puede leer:
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.paramses la solicitud en bruto, así que el propio middleware la valida comoSubscriptionsListenRequestParamsy lee el filtro que pidió el cliente.- El rechazo es un
MCPErrorlanzado antes decall_next(ctx): el cliente recibe ese error y ningún flujo, y la conexión sigue. Mantén el mensaje uniforme, sin nombrar ningún URI, para que un rechazo nunca confirme qué URI están protegidos. - Un único
can_access(user, uri)responde ambas preguntas. El handler del recurso lo consulta enresources/read; el middleware lo consulta ensubscriptions/listen. Cambia la tabla por una base de datos o por tu sistema RBAC y ambos siguen coordinados. - La decisión vale durante toda la vida del flujo. No hay nueva comprobación por evento, así que si el acceso de un llamador puede caducar a mitad del flujo (un token que expira), termina la conexión de ese llamador cuando ocurra.
El contrato completo del middleware, incluido qué más envuelve y por qué está marcado como provisional, está en Middleware.
El lado del cliente
Aquí tienes un cliente al otro lado de ese flujo, siguiendo el tablero:
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)
Entrar en client.listen(...) envía la solicitud y espera tu acuse de recibo, así que el flujo está activo cuando empieza el bloque, y cada evento tipado es una señal para volver a consultar, nunca un payload. Ese es todo el contrato en una pantalla. Todo lo demás sobre el lado del cliente vive en su propia página: observar junto a un flujo principal, finales de flujo y volver a escuchar. Consulta Suscripciones en Clientes.
Escalar más allá de un proceso
Las publicaciones viajan desde tu handler hasta los flujos abiertos a través de un SubscriptionBus. El valor por defecto es en memoria: un proceso, con todos los flujos dentro. Esa es la respuesta correcta hasta que ejecutas réplicas detrás de un balanceador de carga, porque entonces el flujo de un cliente queda fijado a una réplica, y una publicación en otra réplica tiene que llegar hasta él.
Esa pieza te toca implementarla a ti: dos métodos sobre tu backend de pub/sub.
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))
encode es tuyo, y también lo es la tarea lectora de cada réplica que decodifica los mensajes que llegan y llama a cada listener registrado. Los listeners son síncronos, no deben lanzar excepciones y se ejecutan en el bucle de eventos del servidor.
El bus transporta valores ServerEvent tipados, cuatro dataclasses pequeñas, nunca JSON-RPC. El marcado, el filtrado y los ciclos de vida de los flujos se quedan en el SDK, así que una implementación del bus no puede romper el protocolo. Solo puede mover eventos entre procesos.
Para publicar desde fuera de una solicitud, construye el bus tú mismo para conservar la referencia. MCPServer crea uno internamente cuando no pasas nada, y no lo expone.
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
La composición de bajo nivel
Abajo, en el Server de bajo nivel, no hay nada preconectado, y las mismas piezas se ensamblan en tres líneas:
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,
)
- El bus es tuyo, así que publicas en él directamente:
await bus.publish(ResourceUpdated(uri=...)). Ponlo donde tus handlers puedan alcanzarlo: el ámbito del módulo aquí, el lifespan en una app más grande. ListenHandler(bus)es el mismo handler que registraMCPServer, yon_subscriptions_listen=es una ranura de handler común y corriente. Pon tu propio callable en esa ranura para otra semántica, y las obligaciones de la especificación pasan a ti: acusar recibo primero, marcar cada trama con el id de suscripción, no entregar nada fuera del filtro.ListenHandler.close()termina cada flujo abierto de forma ordenada. Cada uno recibe el resultado de la solicitud listen como trama final, que es la forma que tiene la especificación de decir que el servidor terminó la suscripción a propósito. Devuelve antes de que esos flujos acaben de vaciarse, así que dales un momento antes de desmontar el transporte. Sin él, los flujos terminan cuando el cliente se desconecta.
Resumen
- Un cliente se apunta con una solicitud
subscriptions/listen, y la respuesta es el flujo. Atenderla viene integrado. - Publicas con
ctx.notify_*, y el SDK hace el trabajo de marcado, filtrado y ciclo de vida. - Los eventos son señales, no payloads. Ambos extremos vuelven a consultar.
- El lado del cliente es
async with client.listen(...): Suscripciones en Clientes tiene todos los detalles. - En el
Serverde bajo nivel ensamblas tú mismo las mismas piezas: un bus,ListenHandler(bus), la ranuraon_subscriptions_listen. - Escalar horizontalmente significa implementar
SubscriptionBus, dos métodos, y pasarlo comoMCPServer(subscriptions=...).
Ejecutar el servidor que atiende todo esto, detrás de una réplica o de veinte, es Desplegar y escalar.