subscriptions/listen opens a request-scoped stream of change notifications: list changes for tools, prompts, and resources, and updates to specific resources. Subscriptions are not router state. The application supplies the events through a source.

One request may name at most 1,000 URIs in resourceSubscriptions; a longer list is refused with -32602.

Sources

A Snodo.Subscription.Source opens a handle for one listen request, agrees to a subset of the requested filter, blocks in next/2 until it has an event, and closes the handle on cancellation, disconnect, completion, or failure.

next/2 may be called for a handle after close/3 has closed it. When a stream ends, the framework stops the worker that pulls events and closes the handle from a different process, so a pull from that worker can reach the source after the close. Return :closed for a closed handle rather than raise or crash the source process.

The framework:

  • writes the required acknowledgement before pulling the first event;
  • never pulls a second event until the previous one is written;
  • closes the source when the process serving the stream exits for any reason, including when it is killed;
  • stamps the originating request ID on every message;
  • drops core events the client did not ask for;
  • sends a terminal response when the source completes.

Configure the source on the server, and advertise listChanged or subscribe only when one is installed (the runtime refuses to advertise them otherwise):

use Snodo.Server,
  name: "my-server",
  version: "1.0.0",
  subscription_source: {MyApp.ChangeSource, source_options},
  capabilities: %{"tools" => %{"listChanged" => true}}

The hub

Applications that do not need a custom source can supervise Snodo.Subscription.Hub and hand its source to the runtime:

{:ok, hub} = Snodo.Subscription.Hub.start_link(name: MyApp.Hub)
runtime = MyServer.runtime(subscription_source: Snodo.Subscription.Hub.source(hub))

Snodo.Subscription.Hub.publish(hub, Snodo.Subscription.Event.tools_list_changed())

The hub keeps a bounded, filter-aware queue per listener with an explicit overflow policy and delivery statistics. It never detects changes or changes the router; the application publishes when its catalog or resources change.

Across connected nodes

The hub is local to one node. A deployment can forward events from its message bus into a hub on each node while keeping the hub as its subscription source. This example uses OTP :pg and needs no extra package. Start the same :pg scope and one bridge on every directly connected Erlang node that serves subscriptions:

defmodule MyApp.SubscriptionBridge do
  use GenServer

  alias Snodo.Subscription.Hub

  @scope MyApp.SubscriptionGroup
  @group {__MODULE__, :events}

  def start_link(hub), do: GenServer.start_link(__MODULE__, hub, name: __MODULE__)
  def publish(event), do: GenServer.call(__MODULE__, {:publish, event})

  @impl GenServer
  def init(hub) do
    :ok = :pg.join(@scope, @group, self())
    {:ok, hub}
  end

  @impl GenServer
  def handle_call({:publish, event}, _from, hub) do
    case Hub.publish(hub, event) do
      {:ok, local_report} ->
        @scope
        |> :pg.get_members(@group)
        |> Enum.reject(&(&1 == self()))
        |> Enum.each(&send(&1, {:subscription_event, event}))

        {:reply, {:ok, local_report}, hub}

      {:error, _reason} = error ->
        {:reply, error, hub}
    end
  end

  @impl GenServer
  def handle_info({:subscription_event, event}, hub) do
    _result = Hub.publish(hub, event)
    {:noreply, hub}
  end
end

Supervise the scope, local hub, and bridge in that order under a dedicated :rest_for_one supervisor. Restarting the scope or hub then restarts the bridge so it rejoins the group with a live hub. Use the same scope and group names on every node in this deployment:

children = [
  %{id: MyApp.SubscriptionGroup, start: {:pg, :start_link, [MyApp.SubscriptionGroup]}},
  {Snodo.Subscription.Hub, name: MyApp.SubscriptionHub},
  {MyApp.SubscriptionBridge, MyApp.SubscriptionHub}
]

{:ok, _supervisor} = Supervisor.start_link(children, strategy: :rest_for_one)

runtime =
  MyServer.runtime(
    subscription_source: Snodo.Subscription.Hub.source(MyApp.SubscriptionHub)
  )

MyApp.SubscriptionBridge.publish(Snodo.Subscription.Event.tools_list_changed())

The bridge validates an event through its local hub before forwarding it. The returned delivery report describes only that local hub. Remote sends have no acknowledgement. Each receiving hub applies its own listener filters and queue limit. For an existing message bus, replace the :pg send and group membership with a bus subscription and publish into the local hub, or implement Snodo.Subscription.Source directly.

:pg membership converges after nodes connect directly; its view does not relay membership through an intermediary node. A new bridge can miss events until peers learn about it. During a network split, nodes cannot deliver events to peers they cannot see. Neither :pg nor the hub replays missed events after reconnection. A node or local scope restart also loses its in-memory listeners and queued events. Clients should reopen streams and refresh the affected catalogs or resources. For durable delivery or high event volume, use an application-owned bus with its own replay and backpressure policy. The hub's per-listener bound does not bound the bridge process's incoming mailbox.

examples/27_distributed_subscriptions.exs checks the bridge with two local hubs and verifies that each hub still filters events. Run the same bridge and scope on separate directly connected nodes for cross-node delivery.

Transports

Stdio multiplexes streams on one connection and ends one with notifications/cancelled. HTTP answers with text/event-stream, disables proxy buffering, sends keepalive comments, and treats a closed socket as cancellation.

The client

Snodo.Client.listen/3 opens a stream over the direct client, stdio, and HTTP. It returns a Snodo.Client.Subscription once the acknowledgement arrives, with the filter the server accepted, and delivers each event to the owning process as a {:notification, method, params} message on demand, with a bounded buffer in between. See The client.

Extensions

An exact-versioned extension may add filter fields and shape its own events on the same lifecycle. The Tasks package uses this for taskIds and notifications/tasks. See Extensions.

Examples

examples/16_subscriptions.exs (a hand-written source), 18_subscription_hub.exs (the hub), 27_distributed_subscriptions.exs (:pg forwarding), and 17_tasks_subscriptions.exs (Tasks status streams, run from extensions/tasks).