Packages
commanded
1.4.6
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.3
1.4.2
1.4.1
1.4.0
1.4.0-rc.0
1.3.1
1.3.0
1.2.0
1.1.1
1.1.0
1.0.1
1.0.0
1.0.0-rc.1
1.0.0-rc.0
0.19.1
0.19.0
0.18.1
0.18.0
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.0
0.16.0-rc.1
0.16.0-rc.0
0.15.1
0.15.0
0.14.0
0.14.0-rc.0
0.13.0
0.12.0
0.11.0
0.10.0
0.9.0
0.8.5
0.8.4
0.8.3
0.8.1
0.8.0
0.7.1
0.6.2
0.6.1
0.6.0
0.4.0
0.3.1
0.3.0
0.2.1
0.2.0
0.1.0
Use Commanded to build your own Elixir applications following the CQRS/ES pattern.
Current section
Files
Jump to
Current section
Files
test/event_store/support/subscriber.ex
defmodule Commanded.EventStore.Subscriber do
use GenServer
alias Commanded.EventStore.Subscriber
defmodule State do
defstruct [
:subscription_opts,
:event_store,
:event_store_meta,
:owner,
:subscription,
received_events: [],
subscribed?: false
]
end
alias Subscriber.State
def start_link(event_store, event_store_meta, owner, subscription_opts \\ []) do
state = %State{
event_store: event_store,
event_store_meta: event_store_meta,
owner: owner,
subscription_opts: subscription_opts
}
GenServer.start_link(__MODULE__, state)
end
def init(%State{} = state) do
%State{
event_store: event_store,
event_store_meta: event_store_meta,
owner: owner,
subscription_opts: opts
} = state
case event_store.subscribe_to(event_store_meta, :all, "subscriber", self(), :origin, opts) do
{:ok, subscription} ->
state = %State{state | subscription: subscription}
{:ok, state}
{:error, error} ->
send(owner, {:subscribe_error, error, self()})
{:ok, state}
end
end
def ack(subscriber, events),
do: GenServer.call(subscriber, {:ack, events})
def subscribed?(subscriber),
do: GenServer.call(subscriber, :subscribed?)
def received_events(subscriber),
do: GenServer.call(subscriber, :received_events)
def handle_call({:ack, events}, _from, %State{} = state) do
%State{
event_store: event_store,
event_store_meta: event_store_meta,
subscription: subscription
} = state
:ok = event_store.ack_event(event_store_meta, subscription, List.last(events))
{:reply, :ok, state}
end
def handle_call(:subscribed?, _from, %State{} = state) do
%State{subscribed?: subscribed?} = state
{:reply, subscribed?, state}
end
def handle_call(:received_events, _from, %State{} = state) do
%State{received_events: received_events} = state
{:reply, received_events, state}
end
def handle_info({:subscribed, subscription}, %State{subscription: subscription} = state) do
%State{owner: owner} = state
send(owner, {:subscribed, self(), subscription})
{:noreply, %State{state | subscribed?: true}}
end
def handle_info({:events, events}, %State{} = state) do
%State{owner: owner, received_events: received_events} = state
send(owner, {:events, self(), events})
state = %State{state | received_events: received_events ++ events}
{:noreply, state}
end
end