Packages
extreme
0.5.3
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/subscription.ex
defmodule Extreme.Subscription do
use GenServer
require Logger
alias Extreme.Messages, as: ExMsg
def start_link(connection, subscriber, read_params) do
GenServer.start_link __MODULE__, {connection, subscriber, read_params}
end
def start_link(connection, subscriber, stream, resolve_link_tos) do
GenServer.start_link __MODULE__, {connection, subscriber, stream, resolve_link_tos}
end
def init({connection, subscriber, {stream, from_event_number, per_page, resolve_link_tos, require_master}}) do
read_params = %{stream: stream, from_event_number: from_event_number, per_page: per_page,
resolve_link_tos: resolve_link_tos, require_master: require_master}
GenServer.cast self, :read_and_stay_subscribed
{:ok, %{subscriber: subscriber, connection: connection, read_params: read_params, status: :initialized, buffered_messages: [], read_until: -1}}
end
def init({connection, subscriber, stream, resolve_link_tos}) do
read_params = %{stream: stream, resolve_link_tos: resolve_link_tos}
GenServer.cast self, :subscribe
{:ok, %{subscriber: subscriber, connection: connection, read_params: read_params, status: :initialized, buffered_messages: [], read_until: -1}}
end
def handle_cast(:read_and_stay_subscribed, state) do
{:ok, subscription_confirmation} = GenServer.call state.connection, {:subscribe, self, subscribe(state.read_params)}
Logger.debug "Successfully subscribed to stream #{inspect subscription_confirmation}"
GenServer.cast self, :read_events
read_until = subscription_confirmation.last_event_number + 1
{:noreply, %{state | read_until: read_until, status: :reading_events}}
end
def handle_cast(:subscribe, state) do
{:ok, subscription_confirmation} = GenServer.call state.connection, {:subscribe, self, subscribe(state.read_params)}
Logger.debug "Successfully subscribed to stream #{inspect subscription_confirmation}"
{:noreply, %{state | status: :subscribed}}
end
def handle_cast(:read_events, %{read_params: %{from_event_number: from}, read_until: from}=state) do
GenServer.cast self, :push_buffered_messages
{:noreply, %{state|status: :pushing_buffered}}
end
def handle_cast(:read_events, state) do
{read_events, keep_reading} = read_events(state.read_params, state.read_until)
state = case keep_reading do
true -> state
false -> %{state|status: :pushing_buffered}
end
state = Extreme.execute(state.connection, read_events)
|> process_response(state)
{:noreply, state}
end
def handle_cast(:push_buffered_messages, state) do
Enum.each state.buffered_messages, fn e -> send state.subscriber, {:on_event, e} end
{:noreply, %{state|status: :subscribed, buffered_messages: []}}
end
def handle_cast({:ok, %Extreme.Messages.StreamEventAppeared{}=e}, %{status: :subscribed}=state) do
send state.subscriber, {:on_event, e.event}
{:noreply, state}
end
def handle_cast({:ok, %Extreme.Messages.StreamEventAppeared{}=e}, state) do
buffered_messages = state.buffered_messages
|> List.insert_at(-1, e.event)
{:noreply, %{state|buffered_messages: buffered_messages}}
end
def process_response({:ok, %ExMsg.ReadStreamEventsCompleted{}=response}, state) do
Logger.debug "Last read event: #{inspect response.next_event_number - 1}"
push_events {:ok, response}, state.subscriber
send_next_request response, state
end
def process_response({:error, :StreamDeleted, %ExMsg.ReadStreamEventsCompleted{}=response}, state) do
Logger.error "Stream is HARD deleted"
push_events {:extreme, :error, :stream_hard_deleted, state.read_params.stream}, state.subscriber
send_next_request response, state
end
def process_response({:error, :NoStream, %ExMsg.ReadStreamEventsCompleted{}=response}, state) do
Logger.warn "Stream doesn't exist yet"
push_events {:extreme, :warn, :no_stream, state.read_params.stream}, state.subscriber
send_next_request response, state
end
defp push_events({:ok, %ExMsg.ReadStreamEventsCompleted{}=response}, subscriber) do
Enum.each response.events, fn e -> send subscriber, {:on_event, e} end
end
defp push_events({:extreme, _, _, _}=msg, subscriber), do: send(subscriber, msg)
defp send_next_request(_, %{status: :pushing_buffered}=state) do
GenServer.cast self, :push_buffered_messages
state
end
defp send_next_request(%{next_event_number: next_event_number}, state) do
GenServer.cast self, :read_events
%{state|read_params: %{state.read_params|from_event_number: next_event_number}}
end
defp read_events(%{from_event_number: from, per_page: per_page}=params, read_until) when from + per_page < read_until do
result = ExMsg.ReadStreamEvents.new(
event_stream_id: params.stream,
from_event_number: from,
max_count: per_page,
resolve_link_tos: params.resolve_link_tos,
require_master: params.require_master
)
{result, true}
end
defp read_events(params, read_until) do
result = ExMsg.ReadStreamEvents.new(
event_stream_id: params.stream,
from_event_number: params.from_event_number,
max_count: read_until - params.from_event_number,
resolve_link_tos: params.resolve_link_tos,
require_master: params.require_master
)
{result, false}
end
defp subscribe(params) do
ExMsg.SubscribeToStream.new(
event_stream_id: params.stream,
resolve_link_tos: params.resolve_link_tos
)
end
end