Перейти к основному содержимому

Подписки

В MCP 2026-07-28 клиент подписывается на серверные уведомления одним долгоживущим запросом subscriptions/listen с фильтром. На стороне сервера писать обработчик не нужно: neva отвечает на subscriptions/listen сама и рассылает ваши обычные вызовы Context во все потоки, чей фильтр их допускает.

Объявите то, что умеете отправлять​

Принятый фильтр — это запрошенный, суженный до объявленных возможностей, так что объявленное сервером и есть то, на что клиент может подписаться:

use neva::prelude::*;

#[tokio::main]
async fn main() {
App::new()
.with_options(|opt| opt
.with_http(|http| http.bind("127.0.0.1:3000").with_endpoint("/mcp"))
.with_tools(|tools| tools.with_list_changed())
.with_prompts(|prompts| prompts.with_list_changed())
.with_resources(|res| res.with_list_changed().with_subscribe()))
.run()
.await;
}
ВозможностьОткрываетУведомление
tools.listChangedtoolsListChangednotifications/tools/list_changed
prompts.listChangedpromptsListChangednotifications/prompts/list_changed
resources.listChangedresourcesListChangednotifications/resources/list_changed
resources.subscriberesourceSubscriptionsnotifications/resources/updated

Категория, о которой клиент просит, но которую сервер не объявляет, исключается из подтверждения, а не приводит к отказу. Подписка всё равно открывается, и клиент сразу узнаёт, что эти типы никогда не придут.

Ваши обработчики не меняются​

Мутирующие методы Context рассылают уведомления сами — все существующие места вызова продолжают работать, а сервер, у которого раньше не было подписок, теперь их питает:

// Отправляет `notifications/tools/list_changed` каждому потоку, который просил
ctx.tools().add(Tool::new("greet", || async { "hello" })).await?;
let _ = ctx.tools().remove("greet").await?;

// `notifications/prompts/list_changed`
let _ = ctx.prompts().remove("summarize").await?;

// `notifications/resources/list_changed`
ctx.resources().add(Resource::new("res://config", "config")).await?;
let _ = ctx.resources().remove("res://config").await?;

// `notifications/resources/updated` — только потокам, где указан этот URI
ctx.resources().notify_updated("res://config").await?;

Реестр живёт в общем McpOptions, поэтому Context любого выполняющегося запроса достаёт до всех живых потоков — уведомление не заперто внутри того запроса, который его породил.

Уведомления журнала и прогресса не подписочные и сохраняют поведение в области запроса: они идут по потоку ответа того запроса, который их вызвал, — см. Журналирование → Доставка.

notifications/tasks в спецификации является категорией подписки, но в SubscriptionFilter его пока нет, поэтому в сборке по умолчанию Context::task_changed некуда доставлять, а клиенты узнают статус задачи опросом tasks/get.

Кто слушает​

ctx.resources().is_subscribed отвечает по живым потокам, так что можно пропустить дорогую локальную работу, которую всё равно никто не получит:

let resources = ctx.resources();
if resources.is_subscribed(&"res://config".into()) {
// перерисовать снимок, обновить кэш — та самая дорогая часть,
// которую стоит пропустить, если на этом узле никто не слушает
}

// Публикуем в любом случае — фильтры подписок сами всё разошлют.
resources.notify_updated("res://config").await?;
is_subscribed знает только про свой узел

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

Именно поэтому notify_updated не делает предпроверку: он публикует безусловно и оставляет маршрутизацию фильтрам — тем самым, которые этим и занимаются. Используйте is_subscribed, чтобы пропустить работу, но никогда — чтобы решить, слать ли уведомление.

Запуск нескольких экземпляров​

Поток subscriptions/listen — это сокет, который держит ровно один процесс, а транспорт без состояния ни к какому экземпляру клиента не привязывает. Поэтому подписчик и запрос, который меняет сервер, регулярно попадают на разные экземпляры:

клиент --- subscriptions/listen ------------> экземпляр A   (поток здесь)
клиент --- tools/call (меняет список) ------> экземпляр B (ctx.tools().add)
у B подписчиков нет
подписчик A ничего не услышит

Подписчику сказали, что его фильтр принят, поэтому потеря выглядит как «сервер никогда не меняется», а не как сбой доставки.

App::with_notification_bus(..) закрывает этот разрыв: каждый экземпляр публикует то, что произвёл, и доставляет полученное в те потоки, которые держит сам.

use neva::prelude::*;
use neva::app::notification_bus::{BusNotification, NotificationBus};
use neva::shared::Stream;

struct RedisBus { /* … */ }

impl NotificationBus for RedisBus {
async fn publish(&self, notification: BusNotification) {
// передать фоновому соединению
}

fn subscribe(&self) -> impl Stream<Item = BusNotification> + Send + 'static {
// уведомления всех экземпляров, включая собственные
}
}

App::new()
.with_notification_bus(RedisBus { /* … */ })
.with_options(|opt| opt.with_default_http())
.run()
.await;

Таблица подписчиков остаётся локальной для узла по построению: половина каждой записи — это дескриптор сокета на конкретном узле, так что общее хранилище всё равно ничего бы не доставило. Подключаемым сделано именно распространение. neva поставляет трейт и локальное поведение по умолчанию; готовые реализации (Redis pub/sub, NATS, Postgres LISTEN/NOTIFY) живут вне крейта — как и для RequestStateStore.

Контракт​

ПравилоПочему
Без подавления эха — subscribe обязан отдавать и то, что опубликовал этот экземплярЛокальная доставка идёт через тот же поток и только через него. Шина, скрывающая от экземпляра его собственные сообщения, глушит его собственных подписчиков. Redis pub/sub, NATS и tokio::sync::broadcast эхо дают по умолчанию
Достаточно at-most-onceПодписка с переполненным буфером отбрасывает уведомление с предупреждением, а не блокирует запрос, который его породил. Повторная доставка после падения экземпляра ничего не даёт: подписки по спецификации не возобновляемы, а клиент, у которого оборвался поток, шлёт subscriptions/listen заново
publish ожидается внутри запроса-источникаМедленная шина замедляет этот запрос. Предпочитайте реализацию, которая передаёт сообщение фоновому соединению, а не ждёт round trip
Завершившийся поток останавливает доставку навсегдаРеализация, умеющая переподключаться, должна делать это внутри потока, а не завершать его

BusNotification сериализуется как то самое тело уведомления, которое он описывает ({"method": …, "params": …}), так что шина, возящая JSON, может отдать его прямо в serde_json в обе стороны. Он не несёт ничего об экземпляре-источнике или подписке-получателе: принимающий экземпляр сам сверяет его со своими фильтрами и штампует каждую копию идентификатором своего потока.

Без шины ничего не меняется. По умолчанию её нет, уведомления идут прямо подписчикам этого экземпляра, и сервер с одним процессом не платит за существование трейта ни каналом, ни аллокацией, ни задачей.

Три вещи, а не две

Развёртывание из нескольких экземпляров без состояния, обслуживающее подписки, настраивает теперь with_request_state_secret, with_request_state_store и with_notification_bus.

Что идёт по проводу​

--> subscriptions/listen  { "notifications": SubscriptionFilter }
<-- notifications/subscriptions/acknowledged { "notifications": …, "_meta": { subscriptionId } }
<-- notifications/tools/list_changed { "_meta": { subscriptionId } }
…
<-- { "id": …, "result": { "resultType": "complete", "_meta": { subscriptionId } } }

Подтверждение всегда первое сообщение в потоке, и каждое сообщение несёт _meta["io.modelcontextprotocol/subscriptionId"], чтобы клиент, у которого по одному каналу идёт несколько подписок, мог их разделить.

Как завершается подписка​

Что происходитГде применимо
notifications/cancelled для запроса listenstdio
Клиент закрывает потокStreamable HTTP
Закрытие транспортаоба
Остановка сервераоба — после корректного пустого результата

По HTTP notifications/cancelled едет отдельным POST и ничего не доказывает о том, кто открыл поток, поэтому корректный механизм там — закрытие тела ответа, и клиент видит Cancelled, а не финальный результат.

Корректное закрытие при остановке

Спецификация говорит, что сервер, завершающий подписку по собственной инициативе, СЛЕДУЕТ сначала отправить пустой результат, чтобы клиент отличил упорядоченное завершение от оборванного соединения.

Именно поэтому остановка двухфазная. Сигнал завершает подписки и ждёт, пока каждый результат дойдёт до исходящего канала, и только потом разбирается транспорт; run при этом дожидается писателей транспорта, так что результат переживает уроненную следом среду выполнения. App::with_shutdown_drain(..) ограничивает это ожидание (по умолчанию 2 секунды) и полностью пропускается, если ни одной подписки нет, — так что сервер, который ими не пользуется, останавливается не медленнее того, который их не умеет.

Транспорты​

ТранспортКак несётся поток
Streamable HTTPPOST с listen получает ответ text/event-stream, и уведомления идут в его теле. Это третий способ превратить POST в поток — наряду с logLevel и progressToken, — и, в отличие от них, он не требует фичи tracing
stdioСообщения перемежаются с выводом в stdout

Под флагом legacy-spec​

Пара RPC-методов возвращается, и подпиской снова владеет сервер: ctx.resources().subscribe(uri), ctx.resources().unsubscribe(&uri) и resource::commands::{SUBSCRIBE, UNSUBSCRIBE} существуют только под legacy-spec. В сборке по умолчанию уберите вызовы подписки из обработчиков — подпиской теперь владеет клиент, и серверу добавлять нечего.

Обучение на примерах​