An opened, request-scoped MCP subscription.
Applications implement Snodo.Subscription.Source; this module owns filter
narrowing, lifecycle safety, bounded pulling, and protocol wire shaping.
A subscription outlives its request, so it keeps only what it serves. The
source's open/3 callback receives the full request context. The context
stored in the subscription drops request parameters, metadata, client
capabilities, and transport request details. Negotiated extension settings
and accepted filter strings are copied out of the request body, so an open
stream does not keep the body alive.
Summary
Functions
Returns the first required acknowledgement notification.
Closes the application-owned source handle once with an explicit reason.
Returns the final successful JSON-RPC response for graceful completion.
Returns a terminal JSON-RPC error response for a failed open stream.
Shapes one requested source event or drops an event outside the accepted filter.
Opens source for the notifications in requested_filter that the server
supports, and returns the opened subscription.
Starts a monitored, demand-driven source worker owned by owner.
Stops a worker and removes its process monitor.
Types
@type t() :: %Snodo.Subscription{ accepted_filter: map(), context: Snodo.Context.t(), extension_filters: %{optional(String.t()) => map()}, extension_registry: Snodo.Extension.Registry.t(), guard: pid(), handle: term(), id: id(), protocol: module(), source: Snodo.Subscription.Source.Config.t() }
Functions
@spec acknowledgement(t()) :: {:ok, map()} | {:error, Snodo.Error.t()}
Returns the first required acknowledgement notification.
@spec close(t(), Snodo.Subscription.Source.close_reason()) :: :ok
Closes the application-owned source handle once with an explicit reason.
@spec completion(t()) :: {:ok, map()} | {:error, Snodo.Error.t()}
Returns the final successful JSON-RPC response for graceful completion.
Returns a terminal JSON-RPC error response for a failed open stream.
@spec notification(t(), Snodo.Subscription.Event.t()) :: {:ok, map()} | :drop | {:error, Snodo.Error.t()}
Shapes one requested source event or drops an event outside the accepted filter.
@spec open(Snodo.Subscription.Source.Config.t(), map(), Snodo.Context.t()) :: {:ok, t()} | {:error, Snodo.Error.t()}
Opens source for the notifications in requested_filter that the server
supports, and returns the opened subscription.
Serve it with start_worker/2. A source guard starts before the opened
subscription is returned. Transports that supply an owner to open/5 are
protected during the handoff; direct callers can hand ownership to the
worker later.
Starts a monitored, demand-driven source worker owned by owner.
After each continue/1, the owner receives
{:mcp_subscription, worker, outcome} where the outcome is {:ok, event},
:closed, or {:error, reason}.
When open/5 received an owner, the guard watches it during the handoff.
This call hands the guard to owner before starting the pull worker. If the
owner exits, the guard closes the source and the worker stops its puller.
Stops a worker and removes its process monitor.
The worker exits without closing the source handle; the subscription guard
handles close/2 and owner exit.