Подписки
Машинный перевод
Эта страница переведена с английской документации автоматически, и основной версией остаётся английская страница. Если что-то читается неправильно, на странице Переводы объясняется, как об этом сообщить.
Каталог сервера не постоянен. Инструменты появляются во время работы, а содержимое по URI ресурса меняется. Клиент узнаёт об этом через client.listen(...): один запрос subscriptions/listen, ответ на который и есть поток. Он остаётся открытым и несёт те уведомления об изменениях, которые запросил клиент.
Эта страница — о клиентской стороне: как открыть поток, наблюдать за ним рядом с основной логикой и обрабатывать его завершение. Публикация изменений, фильтрация и обслуживание метода — серверная сторона, о ней рассказано на странице Подписки в разделе Внутри обработчика. Примеры здесь общаются с сервером спринт-доски, построенным там.
Наблюдение за потоком
Подписка — это один контекстный менеджер. Вход в него отправляет запрос (именованные аргументы становятся фильтром подписки) и дожидается подтверждения от сервера, так что к началу блока поток уже работает.
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)
Итерация выдаёт четыре типизированных события: ToolsListChanged, PromptsListChanged, ResourcesListChanged и ResourceUpdated(uri=...).
Событие говорит, что изменилось, но никогда — как. Поэтому follow_board вызывает read_resource и list_tools: событие — это сигнал запросить данные заново. Читайте event.uri, а не предполагайте, какой ресурс изменился: фильтр может перечислять несколько URI, а сервер может сообщить об изменении подресурса одного из них.
Дубликаты событий, ожидающих обработки, схлопываются в одно, а повторный запрос всё равно даёт актуальное состояние. Схлопываются только одинаковые события: два ResourceUpdated для разных URI — это два события.
Ещё два свойства дескриптора:
sub.honored— фильтр, который подтвердил сервер:SubscriptionFilterс переданными вами полями, доступными как атрибуты (sub.honored.prompts_list_changed).MCPServerпринимает все виды, которые вы запросили, поэтому возвращает запрос как есть. Сервер, поддерживающий меньше видов, подтверждает меньше, а подтверждённый вид всё равно может ни разу не сработать. Сервер может и отклонить запрос целиком вместо подтверждения (см. Кому разрешено наблюдать на странице сервера) — это проявится как ошибка запроса.sub.subscription_id— идентификатор запроса listen, тот самый, что проставлен на каждом кадре этого потока. Одновременно может быть открыто несколько подписок, и каждая демультиплексируется по своему идентификатору.
Наблюдение без блокировки
follow_board работает, пока сервер не закроет поток, а этого может не случиться никогда, так что сама по себе она забирает всю программу. Настоящим клиентам наблюдатель нужен рядом с основной логикой: агент вызывает инструменты, пока наблюдатель поддерживает актуальность кэша или интерфейса.
Сначала откройте подписку, затем запустите задачу-наблюдатель и занимайтесь своей работой.
import asyncio
from mcp import Client
from mcp.client.subscriptions import Subscription
from .tutorial003 import BOARD, read_board
async def watch(client: Client, sub: Subscription) -> None:
async for _event in sub:
board = await read_board(client)
print(board)
if "[ ]" not in board:
return # sprint finished: the stream closes when run_sprint leaves the block
async def run_sprint(client: Client) -> None:
async with client.listen(resource_subscriptions=[BOARD]) as sub:
print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed
watcher = asyncio.create_task(watch(client, sub))
for task in ("design", "build", "ship"):
await client.call_tool("complete_task", {"board": "sprint", "task": task})
await watcher # returns once the watcher has seen the finished board
async def main() -> None:
async with Client("http://localhost:8000/mcp") as client:
await run_sprint(client)
if __name__ == "__main__":
asyncio.run(main())
import trio
from mcp import Client
from mcp.client.subscriptions import Subscription
from .tutorial003 import BOARD, read_board
async def watch(client: Client, sub: Subscription) -> None:
async for _event in sub:
board = await read_board(client)
print(board)
if "[ ]" not in board:
return # sprint finished: the stream closes when run_sprint leaves the block
async def run_sprint(client: Client) -> None:
async with client.listen(resource_subscriptions=[BOARD]) as sub:
print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed
async with trio.open_nursery() as nursery:
nursery.start_soon(watch, client, sub)
for task in ("design", "build", "ship"):
await client.call_tool("complete_task", {"board": "sprint", "task": task})
async def main() -> None:
async with Client("http://localhost:8000/mcp") as client:
await run_sprint(client)
if __name__ == "__main__":
trio.run(main)
import anyio
from mcp import Client
from mcp.client.subscriptions import Subscription
from .tutorial003 import BOARD, read_board
async def watch(client: Client, sub: Subscription) -> None:
async for _event in sub:
board = await read_board(client)
print(board)
if "[ ]" not in board:
return # sprint finished: the stream closes when run_sprint leaves the block
async def run_sprint(client: Client) -> None:
async with client.listen(resource_subscriptions=[BOARD]) as sub:
print(await read_board(client)) # snapshot: acknowledged, so nothing after this is missed
async with anyio.create_task_group() as tg:
tg.start_soon(watch, client, sub)
for task in ("design", "build", "ship"):
await client.call_tool("complete_task", {"board": "sprint", "task": task})
async def main() -> None:
async with Client("http://localhost:8000/mcp") as client:
await run_sprint(client)
if __name__ == "__main__":
anyio.run(main)
Note
app.py импортирует BOARD и read_board из первого примера, который в этом репозитории
хранится как tutorial003.py. Если вы сохраните показанные файлы рядом под именами client.py
и app.py, напишите вместо этого from client import BOARD, read_board. Пример watch.py
ниже импортирует read_board так же.
Всё дело в порядке. Ничего не воспроизводится повторно, поэтому событие, опубликованное до появления потока, теряется. Вход в client.listen(...) дожидается подтверждения, так что каждое изменение с этого момента доходит до наблюдателя, и снимок, сделанный внутри блока, не может ни одно пропустить.
Запросы свободно выполняются рядом с открытым потоком — из задачи-наблюдателя или любой другой, на том же клиенте. Поскольку одинаковые необработанные события объединяются, загруженная основная логика может дать один повторный запрос вместо трёх. Разные события не объединяются: фильтр со многими URI ставит в очередь по одному ожидающему событию на каждый URI.
Чтобы прекратить наблюдение, выйдите из блока: вызова unsubscribe нет. Отмена задачи, владеющей блоком, делает это за вас, а SDK отменяет запрос listen так, как того ожидает транспорт: в Streamable HTTP — закрытием потока этого запроса. Наблюдатель, работающий всё время жизни приложения, сам никогда не завершится, поэтому отмените его (или область его группы задач) при завершении работы.
Потоки заканчиваются
Поток заканчивается одним из двух способов, и оба — штатный ход выполнения. Корректное закрытие сервером завершает async for; резкий обрыв выбрасывает SubscriptionLost.
Разница диагностическая, а не в том, что делать дальше: потока больше нет, ничего не воспроизводилось повторно, и наблюдатель, которому это всё ещё нужно, слушает заново и перезапрашивает данные.
import anyio
from mcp import Client
from mcp.client.subscriptions import SubscriptionLost
from .tutorial003 import read_board
async def keep_following(client: Client) -> None:
while True:
try:
async with client.listen(resource_subscriptions=["board://sprint"]) as sub:
print(await read_board(client)) # refetch: no replay across streams
async for _event in sub:
print(await read_board(client))
except SubscriptionLost:
pass
# Either ending means the stream is gone. Back off before re-listening:
# a graceful close may be the server shedding load.
await anyio.sleep(1)
Серверы корректно закрывают потоки по своим причинам, в том числе чтобы сбросить подписчика, чья очередь слишком разрослась, поэтому чистое завершение — не сигнал прекратить наблюдение. Перед повторным прослушиванием сделайте паузу.
У SubscriptionLost есть и одна локальная причина. Клиент хранит не более 1024 необработанных событий, и потребитель, отставший настолько, теряет подписку, а не растёт без ограничений. Держите тело async for коротким, а медленную работу выполняйте в другом месте.
keep_following перехватывает только SubscriptionLost. Вход в listen() может также выбросить MCPError (сбой подключения или сервер не обслуживает метод), TimeoutError (подтверждение не пришло) и ListenNotSupportedError (подключение до поколения 2026). Решите, какие из них наблюдателю стоит повторять: последняя не проходит никогда.
Итоги
- Входите в
async with client.listen(...); вход дожидается подтверждения, поэтому ничего из опубликованного после него не теряется. - Итерируйте с помощью
async for event in sub. События — сигналы запросить данные заново, а не сами данные. - Откройте подписку, затем запустите наблюдатель отдельной задачей — и вызовы инструментов продолжают идти рядом.
- Чистое завершение останавливает цикл; обрыв выбрасывает
SubscriptionLost. В обоих случаях: слушайте заново, перезапросите данные, но сначала сделайте паузу. - Выход из блока и есть отписка.
Публикация этих событий, сужение фильтра и масштабирование за пределы одного процесса — серверная сторона: Подписки. Эти же события помогают клиентскому кэшу оставаться актуальным, и следующая страница — Кэширование.