Packages
commanded
0.11.0
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_adapter/subscription_test.exs
defmodule Commanded.EventStore.Adapter.SubscriptionTest do
use Commanded.StorageCase
use Commanded.EventStore
alias Commanded.EventStore.EventData
alias Commanded.Helpers.Wait
defmodule BankAccountOpened, do: defstruct [:account_number, :initial_balance]
defmodule Subscriber do
use GenServer
use Commanded.EventStore
defmodule State do
defstruct [
received_events: [],
subscription: nil,
]
end
alias Subscriber.State
def start_link do
GenServer.start_link(__MODULE__, %State{})
end
def init(%State{} = state) do
{:ok, subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
state = %State{state |
subscription: subscription,
}
{:ok, state}
end
def received_events(subscriber) do
GenServer.call(subscriber, :received_events)
end
def handle_call(:received_events, _from, %State{received_events: received_events} = state) do
{:reply, received_events, state}
end
def handle_info({:events, events}, %State{received_events: received_events, subscription: subscription} = state) do
state = %State{state |
received_events: Enum.concat(received_events, events),
}
@event_store.ack_event(subscription, List.last(events))
{:noreply, state}
end
end
describe "subscribe to all streams" do
test "should receive events appended to any stream" do
{:ok, subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
wait_for_event_store()
{:ok, 1} = @event_store.append_to_stream("stream1", 0, build_events(1))
{:ok, 2} = @event_store.append_to_stream("stream2", 0, build_events(2))
{:ok, 3} = @event_store.append_to_stream("stream3", 0, build_events(3))
assert_receive_events(subscription, 1, from: 1)
assert_receive_events(subscription, 2, from: 2)
assert_receive_events(subscription, 3, from: 4)
refute_receive({:events, _events})
end
test "should skip existing events when subscribing from current position" do
{:ok, 1} = @event_store.append_to_stream("stream1", 0, build_events(1))
{:ok, 2} = @event_store.append_to_stream("stream2", 0, build_events(2))
wait_for_event_store()
{:ok, subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :current)
wait_for_event_store()
refute_receive({:events, _events})
{:ok, 3} = @event_store.append_to_stream("stream3", 0, build_events(3))
wait_for_event_store()
assert_receive_events(subscription, 3, from: 4)
refute_receive({:events, _events})
end
test "should prevent duplicate subscriptions" do
{:ok, _subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
assert {:error, :subscription_already_exists} == @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
end
end
describe "catch-up subscription" do
test "should receive any existing events" do
{:ok, 1} = @event_store.append_to_stream("stream1", 0, build_events(1))
{:ok, 2} = @event_store.append_to_stream("stream2", 0, build_events(2))
{:ok, subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
wait_for_event_store()
assert_receive_events(subscription, 1, from: 1)
assert_receive_events(subscription, 2, from: 2)
{:ok, 3} = @event_store.append_to_stream("stream3", 0, build_events(3))
assert_receive_events(subscription, 3, from: 4)
refute_receive({:events, _events})
end
end
describe "unsubscribe from all streams" do
test "should not receive further events appended to any stream" do
{:ok, subscription} = @event_store.subscribe_to_all_streams("subscriber", self(), :origin)
{:ok, 1} = @event_store.append_to_stream("stream1", 0, build_events(1))
wait_for_event_store()
assert_receive_events(subscription, 1, from: 1)
:ok = @event_store.unsubscribe_from_all_streams("subscriber")
wait_for_event_store()
{:ok, 2} = @event_store.append_to_stream("stream2", 0, build_events(2))
{:ok, 3} = @event_store.append_to_stream("stream3", 0, build_events(3))
refute_receive({:events, _events})
end
end
describe "resume subscription" do
test "should remember last seen event number when subscription resumes" do
{:ok, 1} = @event_store.append_to_stream("stream1", 0, build_events(1))
{:ok, 2} = @event_store.append_to_stream("stream2", 0, build_events(2))
{:ok, subscriber} = Subscriber.start_link()
wait_until(fn ->
received_events = Subscriber.received_events(subscriber)
assert length(received_events) == 3
end)
# wait for last `ack`
:timer.sleep(event_store_wait(200))
Commanded.Helpers.Process.shutdown(subscriber)
wait_for_event_store()
{:ok, subscriber} = Subscriber.start_link()
received_events = Subscriber.received_events(subscriber)
assert length(received_events) == 0
{:ok, 1} = @event_store.append_to_stream("stream3", 0, build_events(1))
wait_until(fn ->
received_events = Subscriber.received_events(subscriber)
assert length(received_events) == 1
end)
end
end
defp wait_until(assertion) do
Wait.until(event_store_wait(1_000), assertion)
end
defp assert_receive_events(subscription, expected_count, opts) do
from_event_number = Keyword.get(opts, :from, 1)
assert_receive {:events, received_events}
received_events
|> Enum.with_index(from_event_number)
|> Enum.each(fn {received_event, expected_event_number} ->
assert received_event.event_number == expected_event_number
end)
@event_store.ack_event(subscription, List.last(received_events))
case expected_count - length(received_events) do
0 -> :ok
remaining when remaining > 0 -> assert_receive_events(subscription, remaining, from: from_event_number + length(received_events))
remaining when remaining < 0 -> flunk("Received #{remaining} more event(s) than expected")
end
end
defp build_event(account_number) do
%EventData{
correlation_id: UUID.uuid4,
event_type: "Elixir.Commanded.EventStore.Adapter.SubscriptionTest.BankAccountOpened",
data: %BankAccountOpened{account_number: account_number, initial_balance: 1_000},
metadata: %{}
}
end
defp build_events(count) do
for account_number <- 1..count, do: build_event(account_number)
end
defp wait_for_event_store do
case event_store_wait() do
nil -> :ok
wait -> :timer.sleep(wait)
end
end
defp event_store_wait(default \\ nil), do: Application.get_env(:commanded, :event_store_wait, default)
end