Snodo.Subscription (snodo v0.4.1)

Copy Markdown View Source

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

id()

@type id() :: integer() | String.t()

t()

@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

acknowledgement(subscription)

@spec acknowledgement(t()) :: {:ok, map()} | {:error, Snodo.Error.t()}

Returns the first required acknowledgement notification.

close(subscription, reason)

@spec close(t(), Snodo.Subscription.Source.close_reason()) :: :ok

Closes the application-owned source handle once with an explicit reason.

completion(subscription)

@spec completion(t()) :: {:ok, map()} | {:error, Snodo.Error.t()}

Returns the final successful JSON-RPC response for graceful completion.

failure(subscription, reason)

@spec failure(t(), term()) :: map()

Returns a terminal JSON-RPC error response for a failed open stream.

notification(subscription, event)

@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.

open(source, requested_filter, context)

@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.

start_worker(subscription, owner)

@spec start_worker(t(), pid()) :: {pid(), reference()}

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.

stop_worker(worker, monitor)

@spec stop_worker(pid(), reference()) :: :ok

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.