Ana içeriğe geç

Abonelikler

Makine çevirisi

Bu sayfa İngilizce dokümantasyondan otomatik olarak çevrildi; esas alınması gereken sürüm İngilizce sayfadır. Yanlış görünen bir şey varsa, nasıl bildireceğinizi Çeviriler sayfası açıklar.

Bir sunucunun kataloğu sabit değildir. Çalışma zamanında yeni araçlar ortaya çıkar, bir kaynak URI'sinin ardındaki içerik değişir.

İstemci bunlardan abonelikler sayesinde haberdar olur. İstemci tek bir subscriptions/listen isteği gönderir ve bu isteğin yanıtı akışın ta kendisidir: açık kalır ve istemcinin istediği değişiklik bildirimlerini taşır.

Değişikliği araçtan yayımlama

Size düşen tek satır: değişikliği yayımlayın.

server.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") bu URI'ye abone olmuş her açık akışa ulaşır. Başka kimseye değil.
  • await ctx.notify_tools_changed() araç listesi değişikliklerini isteyen her akışa ulaşır. Bunu alan istemci tools/list'i yeniden çağırır ve artık sprint_report'u görür.
  • Kardeş metotlar notify_prompts_changed() ve notify_resources_changed().
  • Abone yoksa iş de yok. Boştaki bir sunucuda yayımlamak hiçbir şey yapmaz; bu yüzden kimsenin dinleyip dinlemediğini asla kontrol etmezsiniz. Neyin değiştiğini bildirirsiniz, o kadar.

MCPServer, subscriptions/listen'ı sizin yerinize sunar. Protokol düzeyindeki yükümlülükler (ilk çerçeve olarak onay, akış başına filtreleme, her çerçevede abonelik kimliği) SDK'nın işidir.

Check

Ağ üzerinde, filtresinde board://sprint geçen bir akış complete_task çalıştıktan sonra şöyle görünür:

{"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"}}}

Güncellemenin neyi taşımadığına dikkat edin: panonun kendisini. Her çerçeve, listen isteğinin JSON-RPC kimliğini _meta altında taşır ve bu kimlik abonelik kimliğidir. Onu istemci üretir: Python Client"listen-1" gibi dizeler kullanır; başka istemciler tamsayı kullanabilir.

Yalnızca istenenler

Filtre bir sözleşmedir. Araç listesi değişikliklerini ve tek bir kaynak URI'sini isteyen bir akış bu iki türü alır, başka hiçbir şeyi almaz. Bir prompt değişikliği yayımlarsanız o akış sessiz kalır.

MCPServer kaynak URI'lerini birebir dize olarak eşleştirir; bu yüzden board://sprint URI'sini belirten bir akış board://sprint/tasks/1 hakkında hiçbir şey duymaz. Belirtim, sunucunun abone olunan bir URI'nin alt kaynağındaki değişikliği bildirmesine izin verir; MCPServer bunu hiç yapmaz ama istemciler bunu bekleyecek şekilde yazılmıştır.

Akışın olmadığı iki şey:

  • Bir yeniden oynatma log'u değildir. Kopan bir akış gitmiştir; kimse bağlı değilken yayımlanan olaylar kuyruğa alınmaz. İstemciler yeniden dinler ve yeniden getirir.
  • 2025 yolu değildir. resources/subscribe çağırmış istemcilere ctx.session.send_resource_updated(uri) hizmet verir. notify_* metotları yalnızca subscriptions/listen akışlarına ulaşır.

Kimin izleyebileceğine karar verme

Varsayılan olarak istenen her tür ve URI kabul edilir: her çağıran, yayımladığınız her URI'yi izleyebilir. Okuma işleyicinize hiçbir şey danışmaz, çünkü kimse okumuyordur. files://{name} işleyicinizin geri çevireceği bir çağıran yine de files://payroll.csv üzerinde bir akış açıp onun değiştiğini, hem de ne zaman değiştiğini öğrenebilir. İçeriği asla öğrenemez ve neyin var olduğunu yoklayamaz; çünkü bilinmeyen bir URI de kabul edilir ve yalnızca hiç tetiklenmez. Dar ama gerçek bir açık; bu yüzden çok kiracılı bir sunucudan kullanıcıya özel URI'ler yayımlamadan önce erişimi denetleyin.

Bu denetimi bir middleware (ara katman) üstlenir. subscriptions/listen isteğini SDK onaylamadan önce görür ve çağıran okuyamayacağı bir şey istediğinde isteği reddeder:

server.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 ham istektir; bu yüzden middleware onu SubscriptionsListenRequestParams olarak kendisi doğrular ve istemcinin istediği filtreyi okur.
  • Reddetmek, call_next(ctx)'ten önce fırlatılan bir MCPError demektir: istemci o hatayı alır, akış almaz ve bağlantı devam eder. Mesajı tek tip tutun ve hiçbir URI adı vermeyin; böylece bir ret hangi URI'lerin korunduğunu asla doğrulamaz.
  • Tek bir can_access(user, uri) her iki soruyu da yanıtlar. Kaynak işleyicisi ona resources/read sırasında sorar; middleware ise subscriptions/listen sırasında. Tabloyu bir veritabanıyla ya da RBAC sisteminizle değiştirin, ikisi de uyumlu kalır.
  • Karar akışın ömrü boyunca geçerlidir. Olay başına yeniden denetim yoktur; bu yüzden bir çağıranın erişimi akış ortasında sona erebiliyorsa (süresi dolan bir token gibi), sona erdiğinde o çağıranın bağlantısını kapatın.

Middleware sözleşmesinin tamamı, başka neleri sardığı ve neden geçici (provisional) olarak işaretlendiği de dahil, Middleware sayfasında.

İstemci tarafı

İşte o akışın diğer ucunda, panoyu takip eden bir istemci:

client.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(...)'a girmek isteği gönderir ve sizin onayınızı bekler; yani blok başladığında akış canlıdır ve türü belirli her olay bir yeniden getirme işaretidir, asla bir yük (payload) değildir. Sözleşmenin tamamı tek bir ekranda bu. İstemci tarafıyla ilgili geri kalan her şey kendi sayfasında: ana akışın yanında izleme, akış sonlanmaları ve yeniden dinleme. İstemciler altındaki Abonelikler sayfasına bakın.

Tek sürecin ötesine ölçekleme

Yayımlar, işleyicinizden açık akışlara bir SubscriptionBus üzerinden gider. Varsayılanı bellek içidir: tek süreç, içindeki tüm akışlar. Bir yük dengeleyicinin arkasında replikalar çalıştırana kadar doğru yanıt budur; çünkü o noktada bir istemcinin akışı tek bir replikaya bağlı kalır ve başka bir replikadaki yayımın ona ulaşması gerekir.

Bu birleşim noktasını siz uygularsınız: pub/sub arka ucunuzun üzerinde iki metot.

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 size aittir; her replikada gelen mesajların kodunu çözüp kayıtlı her dinleyiciyi çağıran okuyucu görev de öyle. Dinleyiciler senkrondur, istisna fırlatmamalıdır ve sunucunun olay döngüsünde çalışır.

Veri yolu türü belirli ServerEvent değerleri taşır (dört küçük dataclass), asla JSON-RPC değil. Damgalama, filtreleme ve akış yaşam döngüleri SDK'da kalır; bu yüzden bir veri yolu uygulaması protokolü bozamaz. Yalnızca olayları süreçler arasında taşıyabilir.

Bir isteğin dışından yayımlamak için veri yolunu kendiniz oluşturun ki referansı elinizde olsun. Hiçbir şey geçirmediğinizde MCPServer içeride bir tane kurar ve onu dışarı açmaz.

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

Düşük düzeyli bileşim

Düşük düzeyli Server'da önceden bağlanmış hiçbir şey yoktur; aynı parçalar üç satırda bir araya gelir:

server.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,
)
  • Veri yolu sizindir; bu yüzden doğrudan ona yayımlarsınız: await bus.publish(ResourceUpdated(uri=...)). İşleyicilerinizin erişebileceği bir yere koyun: burada modül kapsamı, daha büyük bir uygulamada lifespan (yaşam döngüsü).
  • ListenHandler(bus), MCPServer'ın kaydettiği işleyicinin aynısıdır ve on_subscriptions_listen= sıradan bir işleyici yuvasıdır. Farklı bir anlam için o yuvaya kendi çağrılabilir nesnenizi koyun; o zaman belirtim yükümlülükleri size geçer: önce onaylayın, her çerçeveyi abonelik kimliğiyle damgalayın, filtrenin dışında hiçbir şey iletmeyin.
  • ListenHandler.close() her açık akışı düzgünce sonlandırır. Her biri son çerçevesi olarak listen isteğinin sonucunu alır; bu, belirtimin sunucunun aboneliği bilerek sonlandırdığını söyleme biçimidir. Metot, bu akışlar boşaltmayı bitirmeden döner; bu yüzden aktarımı kapatmadan önce onlara kısa bir süre tanıyın. Onsuz, akışlar istemci bağlantıyı kestiğinde sona erer.

Özet

  • İstemci tek bir subscriptions/listen isteğiyle katılır ve yanıt akışın kendisidir. Bunu sunmak yerleşiktir.
  • ctx.notify_* ile yayımlarsınız; damgalama, filtreleme ve yaşam döngüsü işini SDK yapar.
  • Olaylar işarettir, yük değil. Her iki uç da yeniden getirir.
  • İstemci tarafı async with client.listen(...) bloğudur: ayrıntıları İstemciler altındaki Abonelikler sayfasında.
  • Düşük düzeyli Server'da aynı parçaları kendiniz birleştirirsiniz: bir veri yolu, ListenHandler(bus), on_subscriptions_listen yuvası.
  • Yatay ölçekleme, SubscriptionBus'ı (iki metot) uygulamak ve onu MCPServer(subscriptions=...) olarak geçirmek demektir.

Tüm bunları sunan sunucuyu ister tek replikanın ister yirmisinin arkasında çalıştırma konusu Dağıtım ve ölçekleme sayfasında.