Подписки
В 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.listChanged | toolsListChanged | notifications/tools/list_changed |
prompts.listChanged | promptsListChanged | notifications/prompts/list_changed |
resources.listChanged | resourcesListChanged | notifications/resources/list_changed |
resources.subscribe | resourceSubscriptions | notifications/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 для запроса listen | stdio |
| Клиент закрывает поток | Streamable HTTP |
| Закрытие транспорта | оба |
| Остановка сервера | оба — после корректного пустого результата |
По HTTP notifications/cancelled едет отдельным POST и ничего не доказывает
о том, кто открыл поток, поэтому корректный механизм там — закрытие тела
ответа, и клиент видит Cancelled, а не финальный результат.
Спецификация говорит, что сервер, завершающий подписку по собственной инициативе, СЛЕДУЕТ сначала отправить пустой результат, чтобы клиент отличил упорядоченное завершение от оборванного соединения.
Именно поэтому остановка двухфазная.
Сигнал завершает подписки и ждёт, пока каждый результат дойдёт до исходящего
канала, и только потом разбирается транспорт; run при этом дожидается
писателей транспорта, так что результат переживает уроненную следом среду
выполнения. App::with_shutdown_drain(..) ограничивает это
ожидание (по умолчанию 2 секунды) и полностью пропускается, если ни одной
подписки нет, — так что сервер, который ими не пользуется, останавливается не
медленнее того, который их не умеет.
Транспорты
| Транспорт | Как несётся поток |
|---|---|
| Streamable HTTP | POST с 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. В сборке по умолчанию уберите вызовы
подписки из обработчиков — подпиской теперь владеет клиент, и серверу
добавлять нечего.
Обучение на примерах
examples/subscriptions— сервер и клиент по HTTPexamples/updates— изменения ресурсов, порождающие уведомления- Клиент → Подписки — вторая половина