Packages

A toolkit for building event-driven applications in Elixir with event sourcing and CQRS patterns

Current section

Files

Jump to
ex_sorcery lib sorcery event_store.ex
Raw

lib/sorcery/event_store.ex

defmodule Sorcery.EventStore do
@moduledoc """
Provides a clean interface for storing and retrieving events.
## Configuration
In your config.exs:
config :sorcery, :event_store,
store: Sorcery.Stores.PostgresStore,
store_opts: [
repo: MyApp.Repo,
table_name: "events"
]
Or for in-memory storage:
config :sorcery, :event_store,
store: Sorcery.Stores.MemoryStore
## Usage
# Create and store an event
event = Sorcery.Event.new(%{
type: "user_registered",
data: %{user_id: "123"},
domain: "users",
instance_id: "instance_1",
domain_sequence_number: 1
})
Sorcery.EventStore.append(event)
# Retrieve events with pagination
{:ok, events, pagination} = Sorcery.EventStore.get_events()
# Query specific events
{:ok, events, pagination} = Sorcery.EventStore.get_events_by_domain("users")
{:ok, events, pagination} = Sorcery.EventStore.get_events_by_type("user_registered")
"""
alias Sorcery.Event
alias Sorcery.EventStore.Behaviour
@type query :: Behaviour.query()
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
end
@doc """
Subscribes to the event store.
"""
@spec subscribe() :: :ok
def subscribe do
Sorcery.PubSub.subscribe()
end
@doc """
Appends one or more events to the configured store.
"""
@spec append(Event.t() | [Event.t()]) :: {:ok, [Event.t()]} | {:error, term()}
def append(events) do
case store().append(events) do
[] ->
{:ok, []}
appended_events ->
Sorcery.PubSub.publish(appended_events)
{:ok, appended_events}
end
end
@doc """
Retrieves events from the configured store.
## 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
## 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()]} | {:error, term()}
def get_events(query \\ nil, opts \\ []) do
store().get_events(query, opts)
end
@doc """
Get events by domain.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `: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
store().get_events_by_domain(domain, opts)
end
@doc """
Get events by type.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `: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
store().get_events_by_type(type, opts)
end
@doc """
Get events by instance.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `: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
store().get_events_by_instance(domain, instance_id, opts)
end
@doc """
Get events in time range.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `: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
store().get_events_in_range(from, to, opts)
end
@doc """
Get events by type in domain.
## Options
* `:limit` - Maximum number of events per page (default: 100)
* `: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
store().get_events_by_type_in_domain(type, domain, opts)
end
@doc """
Returns the configured event store module.
"""
@spec store() :: module()
def store do
Application.get_env(:sorcery, :event_store, [])
|> Keyword.get(:store, Sorcery.EventStore.MemoryStore)
end
@doc """
Returns the configured options for the event store.
"""
@spec store_opts() :: Keyword.t()
def store_opts do
Application.get_env(:sorcery, :event_store, [])
|> Keyword.get(:store_opts, [])
end
@doc """
Starts the configured event store.
"""
@spec start_link(Keyword.t()) :: GenServer.on_start()
def start_link(opts \\ []) do
store().start_link(store_opts() ++ opts)
end
end