Abonnements
Maschinelle Übersetzung
Diese Seite wurde automatisch aus der englischen Dokumentation übersetzt, und die englische Seite ist die maßgebliche Fassung. Wenn sich etwas falsch liest, erklärt Übersetzungen, wie du es melden kannst.
Der Katalog eines Servers ist nicht fest. Tools tauchen zur Laufzeit auf, und der Inhalt hinter einem Ressourcen-URI ändert sich.
Über Abonnements erfährt ein Client davon. Der Client sendet einen einzigen subscriptions/listen-Request, und die Response auf diesen Request ist der Stream: Er bleibt offen und trägt die Änderungsbenachrichtigungen, die der Client angefordert hat.
Aus dem Tool heraus veröffentlichen
Dein Anteil daran ist eine Zeile: Veröffentliche die Änderung.
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")erreicht jeden offenen Stream, der diesen URI abonniert hat. Sonst niemanden.await ctx.notify_tools_changed()erreicht jeden Stream, der Änderungen an der Tool-Liste angefordert hat. Ein Client, der das empfängt, rufttools/listerneut auf und sieht jetztsprint_report.- Die Geschwister heißen
notify_prompts_changed()undnotify_resources_changed(). - Keine Abonnenten, keine Arbeit. Auf einem untätigen Server zu veröffentlichen ist ein No-op, deshalb prüfst du nie, ob jemand zuhört. Du gibst an, was sich geändert hat.
MCPServer bedient subscriptions/listen für dich. Die Pflichten auf der Leitung (die Bestätigung als erster Frame, das Filtern pro Stream, die Abonnement-ID auf jedem Frame) sind Sache des SDK.
Check
Auf der Leitung sieht ein Stream, dessen Filter board://sprint nannte, so aus, nachdem complete_task gelaufen ist:
{"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"}}}
Beachte, was das Update nicht trägt: das Board. Jeder Frame trägt die JSON-RPC-ID des listen-Requests unter _meta, und diese ID ist die Abonnement-ID. Der Client vergibt sie: Der Python-Client verwendet Strings wie "listen-1"; andere Clients verwenden vielleicht Ganzzahlen.
Nur das, was angefordert wurde
Der Filter ist ein Vertrag. Ein Stream, der Änderungen an der Tool-Liste und einen Ressourcen-URI angefordert hat, empfängt diese beiden Arten und nichts anderes. Veröffentlichst du eine Prompt-Änderung, bleibt dieser Stream still.
MCPServer vergleicht Ressourcen-URIs als exakte Strings, deshalb hört ein Stream, der board://sprint nannte, nichts über board://sprint/tasks/1. Die Spezifikation erlaubt einem Server, eine Änderung an einer Unterressource eines abonnierten URI zu melden; MCPServer tut das nie, aber Clients sind darauf ausgelegt, damit zu rechnen.
Zwei Dinge, die der Stream nicht ist:
- Er ist kein Wiederholungsprotokoll. Ein abgebrochener Stream ist weg, und Ereignisse, die veröffentlicht wurden, während niemand verbunden war, werden nicht zwischengespeichert. Clients horchen erneut und laden neu.
- Er ist nicht der Pfad von 2025. Clients, die
resources/subscribeaufgerufen haben, werden überctx.session.send_resource_updated(uri)bedient. Dienotify_*-Methoden erreichen nursubscriptions/listen-Streams.
Entscheiden, wer zusehen darf
Standardmäßig wird jede angeforderte Art und jeder URI akzeptiert: Jeder Aufrufer darf jeden URI beobachten, den du veröffentlichst. Nichts befragt deinen Lese-Handler, weil niemand liest – ein Aufrufer, den dein files://{name}-Handler abweisen würde, kann trotzdem einen Stream auf files://payroll.csv öffnen und erfahren, dass und wann sich die Datei geändert hat. Er erfährt nie Inhalte, und er kann nicht ertasten, was existiert, denn ein unbekannter URI wird ebenfalls akzeptiert und feuert schlicht nie. Schmal, aber real – sichere es also ab, bevor du personenbezogene URIs von einem mandantenfähigen Server veröffentlichst.
Die Absicherung ist eine Middleware. Sie sieht den subscriptions/listen-Request, bevor das SDK ihn bestätigt, und lehnt ab, wenn der Aufrufer etwas anfordert, das er nicht lesen darf:
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.paramsist der rohe Request, deshalb validiert die Middleware ihn selbst zuSubscriptionsListenRequestParamsund liest den Filter, den der Client angefordert hat.- Eine Ablehnung ist ein ausgelöster
MCPErrorvorcall_next(ctx): Der Client bekommt diesen Fehler und keinen Stream, und die Verbindung läuft weiter. Halte die Meldung einheitlich und nenne keinen URI, damit eine Ablehnung nie bestätigt, welche URIs geschützt sind. - Ein einziges
can_access(user, uri)beantwortet beide Fragen. Der Ressourcen-Handler fragt es beiresources/read; die Middleware fragt es beisubscriptions/listen. Tausche die Tabelle gegen eine Datenbank oder dein RBAC-System aus, und beide bleiben im Gleichschritt. - Die Entscheidung gilt für die Lebensdauer des Streams. Es gibt keine erneute Prüfung pro Ereignis. Kann der Zugriff eines Aufrufers also mitten im Stream erlöschen (ein ablaufendes Token), beende die Verbindung dieses Aufrufers, sobald das geschieht.
Der vollständige Middleware-Vertrag, einschließlich dessen, was sie sonst noch umschließt und warum sie als vorläufig markiert ist, steht auf Middleware.
Die Client-Seite
Hier ist ein Client auf der anderen Seite dieses Streams, der dem Board folgt:
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)
Beim Betreten von client.listen(...) wird der Request gesendet und auf deine Bestätigung gewartet, sodass der Stream aktiv ist, wenn der Block beginnt, und jedes typisierte Ereignis ist ein Signal zum Neuladen, nie eine Payload. Das ist der ganze Vertrag auf einem Bildschirm. Alles andere zur Client-Seite steht auf einer eigenen Seite: neben einem Hauptablauf beobachten, Stream-Enden und erneutes Horchen. Siehe Abonnements unter Clients.
Über einen Prozess hinaus skalieren
Veröffentlichungen wandern von deinem Handler über einen SubscriptionBus zu den offenen Streams. Der Standard arbeitet im Speicher: ein Prozess, jeder Stream darin. Das ist die richtige Antwort, bis du Replikate hinter einem Load Balancer betreibst, denn dann ist der Stream eines Clients an ein Replikat gebunden, und eine Veröffentlichung auf einem anderen Replikat muss ihn erreichen.
Diese Nahtstelle implementierst du selbst: zwei Methoden über deinem Pub/Sub-Backend.
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 gehört dir, ebenso der Lese-Task auf jedem Replikat, der eintreffende Nachrichten dekodiert und jeden registrierten Listener aufruft. Listener sind synchron, dürfen keine Exception auslösen und laufen auf der Event-Loop des Servers.
Der Bus trägt typisierte ServerEvent-Werte, vier kleine Dataclasses, nie JSON-RPC. Stempeln, Filtern und Stream-Lebenszyklen bleiben im SDK, sodass eine Bus-Implementierung das Protokoll nicht brechen kann. Sie kann nur Ereignisse zwischen Prozessen bewegen.
Um außerhalb eines Requests zu veröffentlichen, erzeuge den Bus selbst, damit du die Referenz hältst. MCPServer baut intern einen, wenn du nichts übergibst, und legt ihn nicht offen.
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
Die Low-Level-Komposition
Unten auf dem Low-Level-Server ist nichts vorverdrahtet, und dieselben Teile setzen sich in drei Zeilen zusammen:
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,
)
- Der Bus gehört dir, also veröffentlichst du direkt darauf:
await bus.publish(ResourceUpdated(uri=...)). Lege ihn dorthin, wo deine Handler ihn erreichen: hier auf Modulebene, in einer größeren App im Lifespan. ListenHandler(bus)ist derselbe Handler, denMCPServerregistriert, undon_subscriptions_listen=ist ein gewöhnlicher Handler-Slot. Setze dein eigenes Callable in diesen Slot für eine andere Semantik, und die Pflichten aus der Spezifikation gehen auf dich über: zuerst bestätigen, jeden Frame mit der Abonnement-ID stempeln, nichts außerhalb des Filters ausliefern.ListenHandler.close()beendet jeden offenen Stream geordnet. Jeder empfängt das Ergebnis des listen-Requests als letzten Frame – so sagt die Spezifikation, dass der Server das Abonnement absichtlich beendet hat. Die Methode kehrt zurück, bevor diese Streams fertig geleert sind, gib ihnen also einen Moment, bevor du den Transport abbaust. Ohne sie enden Streams, wenn der Client die Verbindung trennt.
Zusammenfassung
- Ein Client steigt mit einem einzigen
subscriptions/listen-Request ein, und die Response ist der Stream. Ihn zu bedienen ist eingebaut. - Du veröffentlichst mit
ctx.notify_*, und das SDK übernimmt Stempeln, Filtern und die Lebenszyklus-Arbeit. - Ereignisse sind Signale, keine Payloads. Beide Seiten laden neu.
- Die Client-Seite ist
async with client.listen(...): Alles Weitere steht in Abonnements unter Clients. - Auf dem Low-Level-
Serversetzt du dieselben Teile selbst zusammen: einen Bus,ListenHandler(bus), den Sloton_subscriptions_listen. - Horizontal skalieren heißt,
SubscriptionBuszu implementieren, zwei Methoden, und ihn alsMCPServer(subscriptions=...)zu übergeben.
Den Server zu betreiben, der all das bedient, hinter einem Replikat oder zwanzig, ist Bereitstellen und skalieren.