Packages
extreme
1.0.2
1.1.4
1.1.3
1.1.2
1.1.1
1.1.1-rc01
1.1.0
1.1.0-rc9
1.1.0-rc8
1.1.0-rc7
1.1.0-rc6
1.1.0-rc5
1.1.0-rc4
1.1.0-rc3
1.1.0-rc2
1.1.0-rc1
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.1
0.12.0
0.11.0
0.10.4
0.10.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
Elixir TCP client for EventStore.
Current section
Files
Jump to
Current section
Files
lib/extreme/subscription.ex
defmodule Extreme.Subscription do
use GenServer
require Logger
alias Extreme.SharedSubscription, as: Shared
defmodule State do
defstruct ~w(base_name correlation_id subscriber stream read_params status)a
end
def start_link(
base_name,
correlation_id,
subscriber,
stream,
resolve_link_tos,
ack_timeout \\ 5_000
) do
GenServer.start_link(
__MODULE__,
{base_name, correlation_id, subscriber, stream, resolve_link_tos, ack_timeout}
)
end
@doc """
Calls `server` with :unsubscribe message. `server` can be either `Subscription` or `ReadingSubscription`.
"""
def unsubscribe(server),
do: GenServer.call(server, :unsubscribe)
@impl true
def init({base_name, correlation_id, subscriber, stream, resolve_link_tos, ack_timeout}) do
read_params = %{stream: stream, resolve_link_tos: resolve_link_tos, ack_timeout: ack_timeout}
state = %State{
base_name: base_name,
correlation_id: correlation_id,
subscriber: subscriber,
read_params: read_params,
status: :initialized
}
{:ok, _} = Shared.subscribe(state)
{:ok, state}
end
@impl true
def handle_call(:unsubscribe, from, state) do
:ok = Shared.unsubscribe(from, state)
{:noreply, state}
end
@impl true
def handle_cast({:process_push, fun}, state),
do: Shared.process_push(fun, state)
end