# Subscriptions

`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):

```elixir
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:

```elixir
{: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`](https://www.erlang.org/docs/27/apps/kernel/pg.html)
and needs no extra package. Start the same `:pg` scope and one bridge on every
directly connected Erlang node that serves subscriptions:

```elixir
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:

```elixir
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](client.md#subscriptions).

## 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](extensions.md).

## 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`).
