Pular para conteúdo

Assinaturas

Tradução automática

Esta página foi traduzida automaticamente a partir da documentação em inglês, e a página em inglês é a versão de referência. Se algo parecer errado, Traduções explica como avisar.

O catálogo de um servidor não é fixo. Ferramentas (tools) aparecem em tempo de execução, e o conteúdo por trás da URI de um recurso muda. Um cliente fica sabendo disso por meio de client.listen(...): uma única requisição subscriptions/listen cuja resposta é o stream. Ele fica aberto e carrega as notificações de mudança que o cliente pediu.

Esta página é a ponta do cliente: abrir o stream, observá-lo ao lado do seu fluxo principal e lidar com seus encerramentos. Publicar mudanças, filtrar e servir o método são o lado do servidor dessa história, contado em Assinaturas, em Dentro do seu handler. Os exemplos aqui conversam com o servidor de quadro de sprint construído lá.

Observando o stream

Uma assinatura é um único gerenciador de contexto. Entrar nele envia a requisição, com seus argumentos nomeados como filtro da assinatura, e espera a confirmação do servidor, então o stream já está ativo quando o bloco começa.

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)

A iteração produz quatro eventos tipados: ToolsListChanged, PromptsListChanged, ResourcesListChanged e ResourceUpdated(uri=...).

Um evento diz o que mudou, nunca como. É por isso que follow_board chama read_resource e list_tools: o evento é uma deixa para buscar de novo. Leia event.uri em vez de presumir qual recurso mudou: um filtro pode nomear várias URIs, e um servidor pode reportar uma mudança em um sub-recurso de uma delas.

Eventos duplicados esperando para serem consumidos se fundem em um só, e buscar de novo ainda traz o estado atual para você. Só eventos idênticos se fundem: dois ResourceUpdated para URIs diferentes são dois eventos.

Mais duas propriedades do handle:

  • sub.honored é o filtro que o servidor confirmou: um SubscriptionFilter com os campos que você passou, lidos como atributos (sub.honored.prompts_list_changed). O MCPServer honra todo tipo que você pede, então ele devolve sua requisição como eco. Um servidor que suporta menos tipos confirma menos, e um tipo honrado ainda pode nunca disparar. Um servidor também pode recusar a requisição inteira em vez de confirmá-la (veja Decidindo quem pode observar na página do servidor), o que aparece como o erro da requisição.
  • sub.subscription_id é o id da requisição listen, aquele carimbado em cada frame deste stream. Várias assinaturas podem estar abertas ao mesmo tempo, cada uma demultiplexada pelo seu próprio id.

Observando sem bloquear

follow_board roda até o servidor fechar o stream, o que pode ser nunca, então sozinha ela toma conta do seu programa. Clientes reais querem o observador ao lado do fluxo principal: um agente chama ferramentas enquanto um observador mantém um cache ou uma UI atualizados.

Abra a assinatura primeiro, depois inicie o observador e siga com o seu trabalho.

app.py
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())
app.py
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)
app.py
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 importa BOARD e read_board do primeiro exemplo, que este repositório guarda como tutorial003.py. Se você salvar os arquivos renderizados lado a lado como client.py e app.py, escreva from client import BOARD, read_board no lugar. O exemplo watch.py mais abaixo importa read_board do mesmo jeito.

A ordem é o ponto. Nada é reenviado, então um evento publicado antes de o seu stream existir se perde. Entrar em client.listen(...) espera a confirmação, então toda mudança daquele momento em diante chega ao seu observador, e o snapshot que você tira dentro do bloco não tem como perder nenhuma.

Requisições rodam livremente ao lado de um stream aberto, a partir da tarefa do observador ou de qualquer outra, no mesmo cliente. Como eventos duplicados não consumidos se fundem, um fluxo principal movimentado pode produzir uma nova busca em vez de três. Eventos diferentes não se fundem: um filtro que nomeia muitas URIs enfileira um evento pendente por URI.

Para parar de observar, saia do bloco: não existe chamada unsubscribe. Cancelar a tarefa que é dona do bloco faz isso por você, e o SDK cancela a requisição listen do jeito que o transporte espera: sobre Streamable HTTP, fechando o stream daquela requisição. Um observador que roda durante toda a vida do seu app nunca retorna sozinho, então cancele-o, ou o escopo do seu task group, no encerramento.

Streams terminam

Um stream termina de uma de duas maneiras, ambas fluxo de controle comum. Um fechamento gracioso do servidor encerra o async for; uma queda abrupta levanta SubscriptionLost.

A diferença é de diagnóstico, não uma diferença no que fazer a seguir: o stream se foi, nada foi reenviado, e um observador que ainda se importa escuta de novo e busca de novo.

watch.py
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)

Servidores fecham streams graciosamente por razões próprias, inclusive para se livrar de um assinante cujo backlog cresceu demais, então um fim limpo não é sinal para parar de observar. Espere um pouco (back off) antes de escutar de novo.

SubscriptionLost também tem uma causa local. O cliente guarda no máximo 1024 eventos não consumidos, e um consumidor que fica tão para trás assim perde a assinatura em vez de crescer sem limite. Mantenha o corpo do async for curto e faça o trabalho lento em outro lugar.

keep_following captura apenas SubscriptionLost. Entrar em listen() também pode levantar MCPError (a conexão falhou, ou o servidor não serve o método), TimeoutError (nenhuma confirmação chegou) e ListenNotSupportedError (uma conexão pré-2026). Decida quais desses o seu observador deve tentar de novo: o último nunca se resolve.

Recapitulando

  • Entre em async with client.listen(...); a entrada espera a confirmação, então nada publicado depois dela se perde.
  • Itere com async for event in sub. Eventos são deixas para buscar de novo, nunca payloads.
  • Abra a assinatura, depois rode o observador como uma tarefa, e as chamadas de ferramentas continuam fluindo ao lado dele.
  • Um fim limpo para o loop; uma queda levanta SubscriptionLost. De qualquer forma: escute de novo, busque de novo, espere um pouco antes.
  • Sair do bloco é o unsubscribe.

Publicar esses eventos, estreitar o filtro e escalar além de um processo são a história do servidor: Assinaturas. Esses mesmos eventos também mantêm um cache do lado do cliente honesto, e Cache é a próxima página.