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
endSupervise 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).