Current section

Files

Jump to
extreme lib subscription.ex
Raw

lib/subscription.ex

defmodule Extreme.Subscription do
use GenServer
require Logger
alias Extreme.Msg, 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)
send state.subscriber, :caught_up
{:noreply, %{state|status: :subscribed, buffered_messages: []}}
end
def handle_cast({:ok, %ExMsg.StreamEventAppeared{}=e}, %{status: :subscribed}=state) do
send(state.subscriber, {:on_event, e.event})
{:noreply, state}
end
def handle_cast({:ok, %ExMsg.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