Packages
x3m_system
0.9.1
0.9.1
0.9.0
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
retired
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
retired
0.7.7
0.7.6
retired
0.7.5
0.7.4
retired
0.7.3
retired
0.7.2
0.7.1
0.7.0
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
retired
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.1.1
0.1.0
Building blocks for distributed and/or CQRS/ES systems
Current section
Files
Jump to
Current section
Files
lib/aggregate_repo.ex
defmodule X3m.System.Aggregate.Repo do
@moduledoc """
Behaviour for the event store that backs aggregates.
An `X3m.System.MessageHandler` reads an aggregate's history and persists its new
events through a module implementing this behaviour, passed as the `:aggregate_repo`
option. The library does not ship a concrete store — you implement these four
callbacks against whatever you use (EventStoreDB/Extreme, Postgres, an in-memory
store for tests, ...):
defmodule MyApp.AggregateRepo do
use X3m.System.Aggregate.Repo
@impl true
def has?(stream_name), do: ...
@impl true
def stream_events(stream_name, start_at, per_page), do: ...
@impl true
def delete_stream(stream_name, hard_delete?, expected_version), do: ...
@impl true
def save_events(stream_name, message, events_metadata), do: ...
end
Stream names are built by the message handler from its `:stream` option and the
aggregate id (`"\#{stream}-\#{id}"`).
"""
@doc """
Returns whether a stream named `stream_name` exists (i.e. the aggregate has any
persisted events).
"""
@callback has?(stream_name :: String.t()) :: boolean
@doc """
Returns an enumerable of `{event, event_number, metadata}` tuples for `stream_name`,
starting at `start_at` and read in pages of `per_page`. Replayed to rebuild state.
"""
@callback stream_events(
stream_name :: String.t(),
start_at :: non_neg_integer(),
per_page :: pos_integer()
) :: Enumerable.t()
@doc """
Deletes `stream_name`. `hard_delete?` chooses a hard vs soft delete; `expected_version`
enables optimistic-concurrency checks (`-2` to skip).
"""
@callback delete_stream(
stream_name :: String.t(),
hard_delete? :: boolean,
expected_version :: integer()
) :: :ok
@doc """
Appends `message.events` to `stream_name`, storing `events_metadata` alongside each
event.
Returns `{:ok, last_event_number}` with the stream's new version on success. Return
`{:error, :wrong_expected_version, expected_last_event_number}` on a concurrency
conflict, or `{:error, reason}` for any other failure (the aggregate process is then
terminated and the error returned to the caller).
"""
@callback save_events(
stream_name :: String.t(),
message :: X3m.System.Message.t(),
events_metadata :: map()
) ::
{:ok, last_event_number :: integer}
| {:error, :wrong_expected_version, expected_last_event_number :: integer}
| {:error, reason :: any}
defmacro __using__(_opts) do
quote do
@moduledoc false
@behaviour X3m.System.Aggregate.Repo
end
end
end