Перейти к содержанию

Подписки

Машинный перевод

Эта страница переведена с английской документации автоматически, и основной версией остаётся английская страница. Если что-то читается неправильно, на странице Переводы объясняется, как об этом сообщить.

Каталог сервера не статичен. Инструменты появляются во время работы, а содержимое, стоящее за URI ресурса, меняется.

Подписки — способ, которым клиент об этом узнаёт. Клиент отправляет один запрос subscriptions/listen, и ответ на этот запрос и есть поток: он остаётся открытым и несёт уведомления об изменениях, которые клиент запросил.

Публикация из инструмента

С вашей стороны нужна одна строка: опубликовать изменение.

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") доходит до каждого открытого потока, подписанного на этот URI. И ни до кого больше.
  • await ctx.notify_tools_changed() доходит до каждого потока, запросившего изменения списка инструментов. Получив его, клиент снова вызывает tools/list и теперь видит sprint_report.
  • Родственные методы — notify_prompts_changed() и notify_resources_changed().
  • Нет подписчиков — нет работы. Публикация на сервере, который никто не слушает, ничего не делает, поэтому проверять, слушает ли кто-нибудь, не нужно. Вы просто сообщаете, что изменилось.

MCPServer обслуживает subscriptions/listen за вас. Протокольные обязательства (подтверждение первым кадром, фильтрация для каждого потока, идентификатор подписки в каждом кадре) — забота SDK.

Check

В передаваемых данных поток, в фильтре которого указан board://sprint, после выполнения 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"}}}

Обратите внимание, чего в обновлении нет: самой доски. Каждый кадр несёт в _meta JSON-RPC-идентификатор запроса listen, и этот идентификатор и есть идентификатор подписки. Его выдаёт клиент: Client на Python использует строки вроде "listen-1", другие клиенты могут использовать целые числа.

Только то, что запрошено

Фильтр — это контракт. Поток, запросивший изменения списка инструментов и один URI ресурса, получает эти два вида событий и ничего больше. Опубликуйте изменение промптов — и этот поток промолчит.

MCPServer сопоставляет URI ресурсов как точные строки, поэтому поток, указавший board://sprint, ничего не услышит о board://sprint/tasks/1. Спецификация разрешает серверу сообщать об изменении подресурса подписанного URI; MCPServer так никогда не делает, но клиенты рассчитаны на такую возможность.

Две вещи, которыми поток не является:

  • Это не журнал для воспроизведения. Оборвавшийся поток потерян, а события, опубликованные, пока никто не был подключён, в очередь не ставятся. Клиенты подписываются заново и заново запрашивают данные.
  • Это не механизм 2025 года. Клиентов, вызвавших resources/subscribe, обслуживает ctx.session.send_resource_updated(uri). Методы notify_* доходят только до потоков subscriptions/listen.

Кто может наблюдать

По умолчанию удовлетворяется каждый запрошенный вид и URI: любой вызывающий может наблюдать за любым URI, который вы публикуете. К вашему обработчику чтения никто не обращается, потому что никто не читает: вызывающий, которому обработчик files://{name} отказал бы, всё равно может открыть поток на files://payroll.csv и узнать, что файл изменился и когда. Содержимого он не узнает никогда и не сможет прощупать, что существует, потому что неизвестный URI тоже принимается и просто никогда не срабатывает. Утечка узкая, но реальная, так что поставьте заслон до того, как публиковать пользовательские URI с мультитенантного сервера.

Заслоном служит middleware (промежуточный слой). Оно видит запрос subscriptions/listen раньше, чем SDK его подтвердит, и отказывает, когда вызывающий просит то, что ему нельзя читать:

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 — это сырой запрос, поэтому middleware само валидирует его в SubscriptionsListenRequestParams и читает фильтр, который запросил клиент.
  • Отказ — это исключение MCPError, выброшенное до call_next(ctx): клиент получает эту ошибку и не получает потока, а соединение продолжает работать. Сообщение делайте единообразным, без упоминания URI, чтобы отказ никогда не подтверждал, какие URI защищены.
  • Одна функция can_access(user, uri) отвечает на оба вопроса. Обработчик ресурса спрашивает её при resources/read, middleware — при subscriptions/listen. Замените таблицу базой данных или своей RBAC-системой, и обе проверки останутся согласованными.
  • Решение действует всё время жизни потока. Повторной проверки на каждое событие нет, поэтому, если доступ вызывающего может истечь посреди потока (токен с ограниченным сроком), завершите его соединение, когда это случится.

Полный контракт middleware, включая то, что ещё оно оборачивает и почему помечено как предварительное, — на странице Middleware.

Клиентская сторона

Вот клиент на другом конце этого потока, следящий за доской:

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(...) отправляет запрос и ждёт вашего подтверждения, так что к началу блока поток уже работает, а каждое типизированное событие — сигнал заново запросить данные, но никогда не сами данные. Вот и весь контракт на одном экране. Всё остальное о клиентской стороне — на отдельной странице: наблюдение параллельно с основным потоком выполнения, завершение потоков и повторная подписка. См. Подписки в разделе Клиенты.

Масштабирование за пределы одного процесса

Публикации идут от обработчика к открытым потокам через SubscriptionBus. По умолчанию шина в памяти: один процесс и все потоки в нём. Это правильный выбор, пока вы не запускаете реплики за балансировщиком нагрузки: тогда поток клиента привязан к одной реплике, а публикация на другой реплике должна до него дойти.

Этот стык реализуете вы: два метода поверх вашего 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 пишете вы, как и задачу-читатель на каждой реплике, которая декодирует приходящие сообщения и вызывает каждого зарегистрированного слушателя. Слушатели синхронны, не должны выбрасывать исключения и выполняются в цикле событий сервера.

Шина несёт типизированные значения ServerEvent — четыре небольших dataclass — и никогда JSON-RPC. Проставление идентификаторов, фильтрация и жизненные циклы потоков остаются в SDK, поэтому реализация шины не может нарушить протокол. Она может лишь переносить события между процессами.

Чтобы публиковать вне запроса, создайте шину сами, чтобы ссылка на неё была у вас. Если ничего не передать, MCPServer создаёт шину внутри и наружу её не отдаёт.

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 ничего заранее не подключено, и те же детали собираются в три строки:

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,
)
  • Шина принадлежит вам, поэтому публикуете вы прямо в неё: await bus.publish(ResourceUpdated(uri=...)). Разместите её там, куда дотянутся обработчики: здесь — на уровне модуля, в приложении побольше — в жизненном цикле (lifespan).
  • ListenHandler(bus) — тот же обработчик, который регистрирует MCPServer, а on_subscriptions_listen= — обычный слот обработчика. Поставьте в этот слот свой вызываемый объект ради другой семантики — и обязательства по спецификации переходят к вам: сначала подтверждение, в каждом кадре идентификатор подписки, ничего за пределами фильтра.
  • ListenHandler.close() корректно завершает все открытые потоки. Каждый получает последним кадром результат запроса listen — так спецификация сообщает, что сервер завершил подписку намеренно. Метод возвращает управление раньше, чем потоки успевают всё отправить, так что дайте им мгновение, прежде чем закрывать транспорт. Без этого вызова потоки заканчиваются, когда отключается клиент.

Итоги

  • Клиент подключается одним запросом subscriptions/listen, и ответом служит поток. Его обслуживание встроено.
  • Вы публикуете через ctx.notify_*, а проставление идентификаторов, фильтрацию и жизненный цикл потоков берёт на себя SDK.
  • События — сигналы, а не данные. Обе стороны запрашивают данные заново.
  • Клиентская сторона — это async with client.listen(...): подробнее — на странице Подписки в разделе Клиенты.
  • На низкоуровневом Server те же детали вы собираете сами: шина, ListenHandler(bus), слот on_subscriptions_listen.
  • Горизонтальное масштабирование — это реализовать SubscriptionBus (два метода) и передать его как MCPServer(subscriptions=...).

О запуске сервера, который всё это обслуживает, с одной репликой или с двадцатью, — на странице Развёртывание и масштабирование.