Snodo.Client.Subscription (snodo v0.4.1)

Copy Markdown View Source

A subscriptions/listen stream opened with Snodo.Client.listen/3.

The struct is the handle: accepted is the filter the server acknowledged, id the request ID (the subscription ID on the wire), ref the tag of the messages the stream sends, owner the process that opened it, and pid the transport process that receives the stream.

Messages

The owner receives {:snodo_subscription, ref, payload} with one of:

  • {:notification, method, params}: one event, as the server sent it. method is the notification's method, for example "notifications/resources/updated", "notifications/tools/list_changed", or an extension's such as "notifications/tasks"; params is its params map, "_meta" included.
  • {:dropped, count}: count events were discarded because the buffer was full. It is sent just before the next delivered event, in addition to it, and takes no demand of its own; that event's unit of demand covers both. A drop happens only while events are queued, so a report always has an event after it, also when the stream is closing.
  • {:closed, reason}: the stream ended. :complete is the server's terminal result; {:error, %Snodo.Error{}} is its terminal error response, or a failure of the connection. Nothing follows it.

Demand and the buffer

Events are sent to the owner only while it has asked for them. demand/2 asks for n more, next/2 asks for one and waits for it, and stream/1 wraps next/2 as an Enumerable. Each unit of demand pays for exactly one event. {:dropped, n} and {:closed, reason} take no demand, so a consumer that renews demand only on events keeps receiving them after a drop. Because next/2 returns one message per call, a call that returns {:dropped, n} leaves the event its demand paid for in the mailbox; the following call returns that event, and the unit of demand that call adds pays for one event delivered ahead of the next call. Events that arrive without demand wait in the transport process, at most :max_buffer of them (100 by default). A full buffer follows the :overflow policy given to Snodo.Client.listen/3: :drop_oldest (the default) discards the oldest queued event, :drop_newest discards the arriving one. Either way the owner is told with {:dropped, n}. {:closed, reason} is delivered after the queued events.

Ending the stream

close/1 ends the stream from the client side; the server sees a cancellation. The owner's exit does the same. No message follows close/1; events delivered before it stay in the owner's mailbox. next/2 on a stream that has ended, after its {:closed, reason} or after close/1, returns {:closed, {:error, %Snodo.Error{}}} at once over every transport.

Summary

Types

Why a stream ended: the server's terminal result, or an error.

What the owner receives inside {:snodo_subscription, ref, payload}.

t()

Functions

Ends the stream.

Asks for n more events, which arrive as messages to the owner.

Asks for one event and waits for the next message.

The stream's payloads as an Enumerable, ending with {:closed, reason}.

Types

close_reason()

@type close_reason() :: :complete | {:error, Snodo.Error.t()}

Why a stream ended: the server's terminal result, or an error.

payload()

@type payload() ::
  {:notification, String.t(), map()}
  | {:dropped, pos_integer()}
  | {:closed, close_reason()}

What the owner receives inside {:snodo_subscription, ref, payload}.

t()

@type t() :: %Snodo.Client.Subscription{
  accepted: map(),
  id: integer(),
  owner: pid(),
  pid: pid(),
  ref: reference()
}

Functions

close(subscription)

@spec close(t()) :: :ok

Ends the stream.

The server is sent a cancellation: stdio writes notifications/cancelled, HTTP closes the connection, and the direct client closes the source. Returns :ok, also for a stream that has already ended.

demand(subscription, n)

@spec demand(t(), pos_integer()) :: :ok

Asks for n more events, which arrive as messages to the owner.

Demand accumulates: queued events are sent at once, up to n, and later events are sent as they arrive until the demand is used up.

next(subscription, timeout \\ :infinity)

@spec next(t(), timeout()) :: payload() | {:error, :timeout}

Asks for one event and waits for the next message.

Returns the payload: {:notification, method, params}, {:dropped, n}, or {:closed, reason}. A {:dropped, n} report is followed by the event this call's demand paid for, which the next call returns. Returns {:error, :timeout} when nothing arrives in timeout milliseconds; the demand stays, so that event arrives as a message later and the next call returns it. If the stream has already ended, or the transport process exits, returns {:closed, {:error, %Snodo.Error{}}}. Must be called by the owner.

stream(subscription)

@spec stream(t()) :: Enumerable.t()

The stream's payloads as an Enumerable, ending with {:closed, reason}.

Each element is fetched with next/2, so the stream runs in the owner. Stopping early leaves the subscription open; call close/1.