Packages
ex_sorcery
0.1.0
A toolkit for building event-driven applications in Elixir with event sourcing and CQRS patterns
Current section
Files
Jump to
Current section
Files
lib/sorcery/event_store/behavior.ex
defmodule Sorcery.EventStore.Behaviour do
@moduledoc """
Defines the behavior that event stores must implement.
An event store is responsible for persisting and retrieving events.
Implementations of this behavior could store events in memory, in a database,
or any other storage mechanism.
## Example Implementation
defmodule MyApp.CustomStore do
use GenServer
@behaviour Sorcery.Stores.Behaviour
# Implement the required callbacks
@impl true
def init(opts) do
{:ok, initial_state}
end
@impl true
def append(state, events) do
# Store the events...
{:ok, new_state}
end
@impl true
def get_events(state, query, opts) do
# Retrieve events based on query...
{:ok, events, new_state}
end
end
## Configuration
config :sorcery, :event_store,
store: MyApp.CustomStore,
store_opts: [...]
"""
alias Sorcery.Event
@type query :: %{
optional(:type) => String.t(),
optional(:domain) => String.t(),
optional(:instance_id) => String.t(),
optional(:from) => DateTime.t(),
optional(:to) => DateTime.t()
}
@doc """
Initializes the event store with the given options.
"""
@callback init_store(opts :: Keyword.t()) :: {:ok, state :: term()}
@doc """
Appends one or more events to the store.
"""
@callback append(state :: term(), events :: Event.t() | [Event.t()]) :: {:ok, [Event.t()], state :: term()}
@doc """
Retrieves events from the store based on query parameters.
## Query Parameters
* `:type` - Filter by event type
* `:domain` - Filter by domain
* `:instance_id` - Filter by instance ID
* `:from` - Filter events after this timestamp
* `:to` - Filter events before this timestamp
The query parameter is optional. Passing nil or an empty map will return all events.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@callback get_events(state :: term(), query :: query() | nil, opts :: Keyword.t()) ::
{:ok, [Event.t()], state :: term()}
defmacro __using__(_opts) do
quote do
use GenServer
@behaviour Sorcery.EventStore.Behaviour
@type query :: %{
optional(:type) => String.t(),
optional(:domain) => String.t(),
optional(:instance_id) => String.t(),
optional(:from) => DateTime.t(),
optional(:to) => DateTime.t()
}
@type pagination :: %{
total_count: non_neg_integer(),
page_size: pos_integer(),
page_number: pos_integer(),
total_pages: pos_integer(),
has_next?: boolean(),
has_prev?: boolean()
}
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
end
@doc """
Starts the store with the given options.
"""
@spec start_link(Keyword.t()) :: GenServer.on_start()
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@impl true
def init(opts) do
init_store(opts)
end
# Helper functions for common queries
@doc """
Get events by domain.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@spec get_events_by_domain(String.t(), Keyword.t()) :: {:ok, [Event.t()]} | {:error, term()}
def get_events_by_domain(domain, opts \\ []) do
GenServer.call(__MODULE__, {:get_events_by_domain, domain, opts})
end
@doc """
Get events by type.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@spec get_events_by_type(String.t(), Keyword.t()) :: {:ok, [Event.t()]} | {:error, term()}
def get_events_by_type(type, opts \\ []) do
GenServer.call(__MODULE__, {:get_events_by_type, type, opts})
end
@doc """
Get events by instance.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@spec get_events_by_instance(String.t(), String.t(), Keyword.t()) :: {:ok, [Event.t()]} | {:error, term()}
def get_events_by_instance(domain, instance_id, opts \\ []) do
GenServer.call(__MODULE__, {:get_events_by_instance, domain, instance_id, opts})
end
@doc """
Get events in time range.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@spec get_events_in_range(DateTime.t(), DateTime.t(), Keyword.t()) :: {:ok, [Event.t()]} | {:error, term()}
def get_events_in_range(from, to, opts \\ []) do
GenServer.call(__MODULE__, {:get_events_in_range, from, to, opts})
end
@doc """
Get events by type in domain.
## Options
* `:limit` - Maximum number of events to return
* `:offset` - Number of events to skip
"""
@spec get_events_by_type_in_domain(String.t(), String.t(), Keyword.t()) :: {:ok, [Event.t()]} | {:error, term()}
def get_events_by_type_in_domain(type, domain, opts \\ []) do
GenServer.call(__MODULE__, {:get_events_by_type_in_domain, type, domain, opts})
end
@doc """
Append events to the store.
"""
@spec append(Event.t() | [Event.t()]) :: {:ok, [Event.t()]} | {:error, term()}
def append(events) do
GenServer.call(__MODULE__, {:append, events})
end
@doc """
Get events from the store.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `:offset` - Number of events to skip
Returns `{:ok, events, pagination}` where pagination includes:
* `:total_count` - Total number of events
* `:page_size` - Number of events per page
* `:page_number` - Current page number
* `:total_pages` - Total number of pages
* `:has_next?` - Whether there are more pages
* `:has_prev?` - Whether there are previous pages
"""
@spec get_events(query() | nil, Keyword.t()) :: {:ok, [Event.t()], pagination()} | {:error, term()}
def get_events(query \\ nil, opts \\ []) do
GenServer.call(__MODULE__, {:get_events, query, opts})
end
# Default handle_call implementations
@impl GenServer
@spec handle_call({:get_events_by_domain, String.t(), Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()]} | {:error, term()}, term()}
def handle_call({:get_events_by_domain, domain, opts}, _from, state) do
handle_get_events(%{domain: domain}, opts, state)
end
@impl GenServer
@spec handle_call({:get_events_by_type, String.t(), Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()]} | {:error, term()}, term()}
def handle_call({:get_events_by_type, type, opts}, _from, state) do
handle_get_events(%{type: type}, opts, state)
end
@impl GenServer
@spec handle_call({:get_events_by_instance, String.t(), String.t(), Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()]} | {:error, term()}, term()}
def handle_call({:get_events_by_instance, domain, instance_id, opts}, _from, state) do
handle_get_events(%{domain: domain, instance_id: instance_id}, opts, state)
end
@impl GenServer
@spec handle_call({:get_events_in_range, DateTime.t(), DateTime.t(), Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()]} | {:error, term()}, term()}
def handle_call({:get_events_in_range, from, to, opts}, _from, state) do
handle_get_events(%{from: from, to: to}, opts, state)
end
@impl GenServer
@spec handle_call({:get_events_by_type_in_domain, String.t(), String.t(), Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()]} | {:error, term()}, term()}
def handle_call({:get_events_by_type_in_domain, type, domain, opts}, _from, state) do
handle_get_events(%{type: type, domain: domain}, opts, state)
end
@impl GenServer
@spec handle_call({:append, Event.t() | [Event.t()]}, GenServer.from(), term()) ::
{:reply, Event.t() | [Event.t()], term()}
def handle_call({:append, events}, _from, state) do
case append(state, events) do
{:ok, appended_events, new_state} -> {:reply, appended_events, new_state}
end
end
@impl GenServer
@spec handle_call({:get_events, query() | nil, Keyword.t()}, GenServer.from(), term()) ::
{:reply, {:ok, [Event.t()], pagination()}, term()}
def handle_call({:get_events, query, opts}, _from, state) do
case validate_query(query) do
:ok -> handle_get_events(query, opts, state)
error -> {:reply, error, state}
end
end
# Private helper for handling get_events calls
@spec handle_get_events(query() | nil, Keyword.t(), term()) ::
{:reply, {:ok, [Event.t()]}, term()}
defp handle_get_events(query, opts, state) do
{:ok, events, new_state} = get_events(state, query, opts)
{:reply, {:ok, events}, new_state}
end
# Filter functions
@doc """
Filters events based on query parameters.
Returns unfiltered events if query is nil or empty.
"""
@spec filter_events([Event.t()], query() | nil) :: [Event.t()]
def filter_events(events, nil), do: events
def filter_events(events, query) when map_size(query) == 0, do: events
def filter_events(events, query) do
events
|> filter_by_type(query)
|> filter_by_domain(query)
|> filter_by_instance_id(query)
|> filter_by_time_range(query)
end
@doc false
@spec filter_by_type([Event.t()], query()) :: [Event.t()]
defp filter_by_type(events, %{type: type}) when not is_nil(type) do
Enum.filter(events, & &1.type == type)
end
defp filter_by_type(events, _), do: events
@doc false
@spec filter_by_domain([Event.t()], query()) :: [Event.t()]
defp filter_by_domain(events, %{domain: domain}) when not is_nil(domain) do
Enum.filter(events, & get_in(&1.metadata, [:domain]) == domain)
end
defp filter_by_domain(events, _), do: events
@doc false
@spec filter_by_instance_id([Event.t()], query()) :: [Event.t()]
defp filter_by_instance_id(events, %{instance_id: instance_id}) when not is_nil(instance_id) do
Enum.filter(events, & get_in(&1.metadata, [:instance_id]) == instance_id)
end
defp filter_by_instance_id(events, _), do: events
@doc false
@spec filter_by_time_range([Event.t()], query()) :: [Event.t()]
defp filter_by_time_range(events, %{from: from, to: to})
when not is_nil(from) and not is_nil(to) do
Enum.filter(events, fn event ->
DateTime.compare(event.inserted_at, from) in [:gt, :eq] and
DateTime.compare(event.inserted_at, to) in [:lt, :eq]
end)
end
defp filter_by_time_range(events, _), do: events
# Query validation
@valid_query_keys [:type, :domain, :instance_id, :from, :to]
@doc """
Validates a query map.
"""
@spec validate_query(query() | nil) :: :ok | {:error, String.t()}
def validate_query(nil), do: :ok
def validate_query(query) when not is_map(query), do: {:error, "Query must be a map or nil"}
def validate_query(query) do
with :ok <- validate_query_keys(query),
:ok <- validate_time_range(query) do
:ok
end
end
@spec validate_query_keys(query()) :: :ok | {:error, String.t()}
defp validate_query_keys(query) do
invalid_keys = Map.keys(query) -- @valid_query_keys
if invalid_keys == [] do
:ok
else
{:error, "Invalid query keys: #{inspect(invalid_keys)}"}
end
end
@spec validate_time_range(query()) :: :ok | {:error, String.t()}
defp validate_time_range(%{from: from, to: to})
when not is_nil(from) and not is_nil(to) do
case DateTime.compare(from, to) do
:gt -> {:error, "From date must be before to date"}
_ -> :ok
end
end
defp validate_time_range(_), do: :ok
# Pagination
@doc """
Applies pagination to a list of events.
"""
@spec paginate([Event.t()], Keyword.t()) :: {[Event.t()], pagination()}
def paginate(events, opts) do
page_size = Keyword.get(opts, :limit, 100)
page_number = div(Keyword.get(opts, :offset, 0), page_size) + 1
total_count = length(events)
total_pages = ceil(total_count / page_size)
paginated_events =
events
|> Enum.drop((page_number - 1) * page_size)
|> Enum.take(page_size)
pagination = %{
total_count: total_count,
page_size: page_size,
page_number: page_number,
total_pages: total_pages,
has_next?: page_number < total_pages,
has_prev?: page_number > 1
}
{paginated_events, pagination}
end
defoverridable []
end
end
end