Packages
ex_esdb
0.0.3-alpha
0.11.0
0.10.0
0.9.0
0.8.0
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.1
0.6.0
0.5.1
0.5.0
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.3.3
0.3.2
0.3.1
0.3.0
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14-alpha
0.0.13-alpha
0.0.12-alpha
0.0.11-alpha
0.0.10-alpha
0.0.9-alpha
0.0.8-alpha
0.0.6-alpha
0.0.5-alpha
0.0.4-alpha
0.0.3-alpha
0.0.2-alfa
0.0.1-alfa
ExESDB is a reincarnation of rabbitmq/khepri, specialized for use as a BEAM-native event store.
Current section
Files
Jump to
Current section
Files
lib/ex_esdb/event_store.ex
defmodule ExESDB.EventStore do
@moduledoc """
A GenServer wrapper around :khepri to act as an event store.
Inspired by EventStoreDB's API.
"""
use GenServer
require Logger
alias ExESDB.EventEmitter, as: ESEmitter
alias ExESDB.EventStreamReader, as: ESReader
alias ExESDB.EventStreamWriter, as: ESWriter
# Client API
@doc """
Get the list of streams in the store.
## Parameters
- `store`: The store to get the streams from.
## Returns
- `{:ok, streams}` if successful.
"""
def get_streams(store),
do:
GenServer.call(
__MODULE__,
{:get_streams, store}
)
@doc """
Get the current state of the store.
## Parameters
- `store`: The store to get the state of.
## Returns
- `{:ok, state}` if successful.
- `{:error, reason}` if unsuccessful.
"""
def get_state(_store),
do:
GenServer.call(
__MODULE__,
{:get_state}
)
@doc """
Append events to a stream.
## Parameters
- `stream_id`: The id of the stream to append to.
- `expected_version`: The expected version of the stream (for optimistic concurrency).
- `events`: A list of events to append.
## Returns
- `{:ok, new_stream_version}` if successful.
- `{:error, reason}` if unsuccessful.
"""
def append_to_stream(store, stream_id, expected_version, events),
do:
GenServer.call(
__MODULE__,
{:append_to_stream, store, stream_id, expected_version, events}
)
@doc """
Read events from a stream.
## Parameters
- `stream_name`: The name of the stream to read from.
- `start_version`: The version to start reading from.
- `count`: The number of events to read.
## Returns
- `{:ok, events}` if successful.
- `{:error, reason}` if unsuccessful.
"""
def read_stream_forward(store, stream_id, start_version, count),
do:
GenServer.call(
__MODULE__,
{:read_stream_forward, store, stream_id, start_version, count}
)
@doc """
Get the current version of a stream.
## Parameters
- `stream_name`: The name of the stream to check.
## Returns
- `{:ok, version}` if successful.
- `{:error, reason}` if unsuccessful.
"""
def stream_version(store, stream_id) do
GenServer.call(
__MODULE__,
{:stream_version, store, stream_id}
)
end
@impl true
def handle_info(:sigterm, %{store_id: store}) do
Logger.info("#{Colors.store_theme(self())} => SIGTERM received. Resetting cluster #{inspect(store, pretty: true)}")
{:stop, :normal, nil}
end
@impl true
def handle_info(:register_emitter, [config: _, store: store] = state) do
store
|> ESEmitter.register_erl_emitter()
{:noreply, state}
end
## CALLBACKS
@impl true
def handle_call({:get_state}, _from, state) do
{:reply, {:ok, state}, state}
end
@impl true
def handle_call(
{:get_streams, store},
_from,
state
) do
streams =
store
|> ESReader.get_streams()
{:reply, {:ok, streams}, state}
end
@impl true
def handle_call(
{:append_to_stream, store, stream_id, expected_version, events},
_from,
state
) do
current_version =
store
|> ESReader.get_current_version!(stream_id)
if current_version == expected_version do
new_version =
store
|> ESWriter.append_events(stream_id, events, current_version)
{:reply, {:ok, new_version}, state}
else
{:reply, {:error, :wrong_expected_version}, state}
end
end
@impl true
def handle_call(
{:read_stream_forward, store, stream_id, start_version, count},
_from,
state
) do
events =
store
|> ESReader.read_events(stream_id, start_version, count)
{:reply, {:ok, events}, state}
end
@impl true
def handle_call(
{:stream_version, store, stream_id},
_from,
state
) do
version =
store
|> ESReader.get_current_version!(stream_id)
{:reply, {:ok, version}, state}
end
defp start_khepri(opts) do
store = opts[:store_id]
timeout = opts[:timeout]
data_dir = opts[:data_dir]
case :khepri.start(data_dir, store, timeout) do
{:ok, store} ->
Logger.info("#{Colors.store_theme(self())} => Started store: #{inspect(store, pretty: true)}")
{:ok, store}
reason ->
Logger.error("Failed to start khepri. reason: #{inspect(reason)}")
end
end
#### PLUMBING
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: 10_000,
type: :worker
}
end
def start_link(opts),
do:
GenServer.start_link(
__MODULE__,
opts,
name: __MODULE__
)
# Server Callbacks
@impl true
def init(opts) do
IO.puts("Starting ExESDB.EventStore with config: #{inspect(opts, pretty: true)}")
Process.flag(:trap_exit, true)
Process.send_after(self(), :register_emitter, 10_000)
case start_khepri(opts) do
{:ok, store} ->
Logger.debug("Started store: #{inspect(store)}")
:os.set_signal(:sigterm, :handle)
{:ok, [config: opts, store: store]}
reason ->
Logger.error("Failed to start khepri. reason: #{inspect(reason)}")
{:error, [config: opts, store: nil]}
end
end
end