Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
86 changes: 16 additions & 70 deletions test/aggregates/aggregate_concurrency_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -3,30 +3,11 @@ defmodule Commanded.Aggregates.AggregateConcurrencyTest do

alias Commanded.MockedApp
alias Commanded.Aggregates.{Aggregate, ExecutionContext}
alias Commanded.EventStore.RecordedEvent
alias Commanded.ExampleDomain.{BankAccount, DepositMoneyHandler, OpenAccountHandler}
alias Commanded.ExampleDomain.BankAccount.Commands.{DepositMoney, OpenAccount}
alias Commanded.ExampleDomain.BankAccount.Events.MoneyDeposited
alias Commanded.UUID

setup do
expect(MockEventStore, :subscribe_to, fn
_event_store_meta, stream_uuid, handler_name, handler, _subscribe_from, _opts ->
assert is_binary(stream_uuid)
assert is_binary(handler_name)

{:ok, handler}
end)

expect(MockEventStore, :subscribe, fn _event_store_meta, aggregate_uuid ->
assert is_binary(aggregate_uuid)

:ok
end)

:ok
end

describe "concurrency error" do
setup [:open_account]

Expand All @@ -46,38 +27,20 @@ defmodule Commanded.Aggregates.AggregateConcurrencyTest do
retry_attempts: 1
}

# Fail to append once
expect(MockEventStore, :append_to_stream, fn
_event_store_meta, ^account_number, 1, _event_data, _opts ->
{:error, :wrong_expected_version}
end)

# Return "missing" event
expect(MockEventStore, :stream_forward, fn
_event_store_meta, ^account_number, 2, _batch_size ->
[
%RecordedEvent{
event_id: UUID.uuid4(),
event_number: 2,
stream_id: account_number,
stream_version: 2,
event_type: "Elixir.Commanded.ExampleDomain.BankAccount.Events.MoneyDeposited",
data: %MoneyDeposited{
account_number: account_number,
transfer_uuid: UUID.uuid4(),
amount: 500,
balance: 1_500
},
metadata: %{}
}
]
end)

# Succeed on second attempt
expect(MockEventStore, :append_to_stream, fn
_event_store_meta, ^account_number, 2, _event_data, _opts ->
:ok
end)
concurrent_event =
build_recorded_event(
account_number,
2,
%MoneyDeposited{
account_number: account_number,
transfer_uuid: UUID.uuid4(),
amount: 500,
balance: 1_500
},
event_type: "Elixir.Commanded.ExampleDomain.BankAccount.Events.MoneyDeposited"
)

expect_concurrency_retry_succeeds(account_number, concurrent_event)

assert {:ok, 3, _events, _aggregate_state} =
Aggregate.execute(MockedApp, BankAccount, account_number, context)
Expand All @@ -92,16 +55,7 @@ defmodule Commanded.Aggregates.AggregateConcurrencyTest do
end

test "should error after too many attempts", %{account_number: account_number} do
# Fail to append to stream
expect(MockEventStore, :append_to_stream, 6, fn
_event_store_meta, ^account_number, 1, _event_data, _opts ->
{:error, :wrong_expected_version}
end)

expect(MockEventStore, :stream_forward, 6, fn
_event_store_meta, ^account_number, 2, _batch_size ->
[]
end)
expect_too_many_retry_attempts(account_number)

command = %DepositMoney{
account_number: account_number,
Expand All @@ -123,15 +77,7 @@ defmodule Commanded.Aggregates.AggregateConcurrencyTest do
defp open_account(_context) do
account_number = UUID.uuid4()

expect(MockEventStore, :stream_forward, fn
_event_store_meta, ^account_number, 1, _batch_size ->
[]
end)

expect(MockEventStore, :append_to_stream, fn
_event_store_meta, ^account_number, 0, _event_data, _opts ->
:ok
end)
expect_open_aggregate(account_number)

{:ok, ^account_number} =
Commanded.Aggregates.Supervisor.open_aggregate(MockedApp, BankAccount, account_number)
Expand Down
45 changes: 8 additions & 37 deletions test/aggregates/aggregate_telemetry_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -258,22 +258,7 @@ defmodule Commanded.Aggregates.AggregateTelemetryTest do
test "emit `[:commanded, :aggregate, :execute, :wrong_expected_version]` event" do
aggregate_uuid = UUID.uuid4()

expect(MockEventStore, :subscribe, fn _event_store_meta, aggregate_uuid ->
assert is_binary(aggregate_uuid)
:ok
end)

expect(MockEventStore, :append_to_stream, fn _meta,
^aggregate_uuid,
_exp_ver,
_event_data,
_opts ->
{:error, :wrong_expected_version}
end)

expect(MockEventStore, :stream_forward, 2, fn _meta, ^aggregate_uuid, _from, _batch_size ->
[]
end)
expect_wrong_expected_version_conflict(aggregate_uuid)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: Commanded.MockedApp)

Expand Down Expand Up @@ -320,11 +305,7 @@ defmodule Commanded.Aggregates.AggregateTelemetryTest do
test "load only when stream_forward returns stream_not_found" do
aggregate_uuid = UUID.uuid4()

expect(MockEventStore, :subscribe, fn _meta, ^aggregate_uuid -> :ok end)

expect(MockEventStore, :stream_forward, fn _meta, ^aggregate_uuid, _from, _batch_size ->
{:error, :stream_not_found}
end)
expect_stream_not_found(aggregate_uuid)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: MockedApp)

Expand All @@ -339,24 +320,14 @@ defmodule Commanded.Aggregates.AggregateTelemetryTest do
aggregate_uuid = UUID.uuid4()
count = 2

expect(MockEventStore, :subscribe, fn _meta, ^aggregate_uuid -> :ok end)

expect(MockEventStore, :stream_forward, fn _meta, ^aggregate_uuid, _from, _batch_size ->
events =
for i <- 1..count do
%Commanded.EventStore.RecordedEvent{
event_id: UUID.uuid4(),
event_number: i,
stream_id: aggregate_uuid,
stream_version: i,
correlation_id: nil,
causation_id: nil,
event_type: "Elixir.Commanded.Aggregates.AggregateTelemetryTest.Event",
data: %Event{message: "event#{i}"},
metadata: nil,
created_at: DateTime.utc_now()
}
build_recorded_event(aggregate_uuid, i, %Event{message: "event#{i}"},
event_type: "Elixir.Commanded.Aggregates.AggregateTelemetryTest.Event"
)
end
end)

expect_stream_with_events(aggregate_uuid, events)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: MockedApp)

Expand Down
9 changes: 1 addition & 8 deletions test/event_handler_after_start_test.exs
Original file line number Diff line number Diff line change
@@ -1,15 +1,8 @@
defmodule Commanded.Event.HandlerAfterStartTest do
use Commanded.MockEventStoreCase

import Mox

alias Commanded.EventStore.Adapters.Mock, as: MockEventStore

setup do
stub(MockEventStore, :subscribe_to, fn
_event_store, :all, _handler_name, handler, _subscribe_from, _opts ->
{:ok, handler}
end)
stub_subscribe_to_return_handler()

:ok
end
Expand Down
79 changes: 10 additions & 69 deletions test/opentelemetry/aggregate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
alias Commanded.Aggregates.AggregateTelemetryTest
alias Commanded.Aggregates.{Aggregate, ExecutionContext}
alias Commanded.DefaultApp
alias Commanded.EventStore.Adapters.Mock, as: MockEventStore
alias Commanded.Middleware.Commands.IncrementCount
alias Commanded.Middleware.Commands.RaiseError
alias Commanded.MockedApp
Expand Down Expand Up @@ -606,22 +605,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
test "span has wrong_expected_version event when conflict occurs" do
aggregate_uuid = UUID.uuid4()

expect(MockEventStore, :subscribe, fn _event_store_meta, ^aggregate_uuid ->
assert is_binary(aggregate_uuid)
:ok
end)

expect(MockEventStore, :append_to_stream, fn _meta,
^aggregate_uuid,
_exp_ver,
_event_data,
_opts ->
{:error, :wrong_expected_version}
end)

expect(MockEventStore, :stream_forward, 2, fn _meta, ^aggregate_uuid, _from, _batch_size ->
[]
end)
expect_wrong_expected_version_conflict(aggregate_uuid)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: MockedApp)

Expand Down Expand Up @@ -673,44 +657,15 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
test "span has wrong_expected_version event with count when retry succeeds" do
aggregate_uuid = UUID.uuid4()

expect(MockEventStore, :subscribe, fn _event_store_meta, ^aggregate_uuid ->
assert is_binary(aggregate_uuid)
:ok
end)

expect(MockEventStore, :append_to_stream, 2, fn
_meta, ^aggregate_uuid, 0, _event_data, _opts ->
{:error, :wrong_expected_version}

_meta, ^aggregate_uuid, 1, _event_data, _opts ->
:ok
end)
event =
build_recorded_event(
aggregate_uuid,
1,
struct!(Commanded.Aggregates.AggregateTelemetryTest.Event, message: "event"),
event_type: "Elixir.Commanded.Aggregates.AggregateTelemetryTest.Event"
)

stream_forward_calls = :counters.new(1, [])

expect(MockEventStore, :stream_forward, 2, fn _meta, ^aggregate_uuid, 1, _batch_size ->
n = :counters.get(stream_forward_calls, 1)
:counters.add(stream_forward_calls, 1, 1)

if n == 0 do
[]
else
[
%Commanded.EventStore.RecordedEvent{
event_id: UUID.uuid4(),
event_number: 1,
stream_id: aggregate_uuid,
stream_version: 1,
correlation_id: nil,
causation_id: nil,
event_type: "Elixir.Commanded.Aggregates.AggregateTelemetryTest.Event",
data: struct!(Commanded.Aggregates.AggregateTelemetryTest.Event, message: "event"),
metadata: nil,
created_at: DateTime.utc_now()
}
]
end
end)
expect_wrong_expected_version_retry_succeeds(aggregate_uuid, event)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: MockedApp)

Expand Down Expand Up @@ -760,21 +715,7 @@ defmodule Commanded.OpenTelemetry.AggregateTest do
test "no wrong_expected_version event when count is 0" do
aggregate_uuid = UUID.uuid4()

expect(MockEventStore, :subscribe, fn _event_store_meta, ^aggregate_uuid ->
:ok
end)

expect(MockEventStore, :append_to_stream, fn _meta,
^aggregate_uuid,
_exp_ver,
_event_data,
_opts ->
:ok
end)

expect(MockEventStore, :stream_forward, fn _meta, ^aggregate_uuid, _from, _batch_size ->
[]
end)
expect_successful_append_with_empty_stream(aggregate_uuid)

assert {:ok, _pid} = start_aggregate(aggregate_uuid, application: MockedApp)

Expand Down
26 changes: 1 addition & 25 deletions test/projections/error_callback_test.exs
Original file line number Diff line number Diff line change
@@ -1,11 +1,9 @@
defmodule Commanded.Projections.ErrorCallbackTest do
use ExUnit.Case
use Commanded.MockProjectionCase

import Commanded.Projections.ProjectionAssertions
import ExUnit.CaptureLog
import Mox

alias Commanded.EventStore.Adapters.Mock, as: MockEventStore
alias Commanded.EventStore.RecordedEvent

alias Commanded.Projections.Events.{
Expand All @@ -21,13 +19,6 @@ defmodule Commanded.Projections.ErrorCallbackTest do
alias Commanded.UUID
alias Ecto.Adapters.SQL.Sandbox

setup [:set_mox_global, :stub_event_store, :verify_on_exit!]

setup do
start_supervised!({TestApplication, event_store: [adapter: MockEventStore]})
Sandbox.checkout(Repo)
end

describe "error handling" do
setup [:start_projector]

Expand Down Expand Up @@ -147,21 +138,6 @@ defmodule Commanded.Projections.ErrorCallbackTest do
end
end

defp stub_event_store(_context) do
stub(MockEventStore, :ack_event, fn _adapter_meta, _pid, _event -> :ok end)

stub(MockEventStore, :child_spec, fn _application, _config ->
{:ok, [], %{}}
end)

stub(MockEventStore, :subscribe_to, fn
_event_store, :all, _handler_name, _handler, _subscribe_from, _opts ->
{:ok, self()}
end)

:ok
end

defp start_projector(_context) do
projector = start_supervised!(ErrorProjector)

Expand Down
Loading
Loading