Packages
extreme
0.6.0
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
state.buffered_messages |> Enum.each(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
response.events |> Enum.each(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