Packages
ex_esdb
0.0.16
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/system.ex
defmodule ExESDB.System do
@moduledoc """
This module is the top level supervisor for the ExESDB system.
It is responsible for supervising:
- The PubSub mechanism
- the Event Store (starts and stops khepri)
- the Cluster (joins and leaves the cluster)
- the Leader (manages Ra leader-specific functionality)
- the Subscriptions Supervisor (manages subscriptions)
"""
use Supervisor
alias ExESDB.Options, as: Options
alias ExESDB.Themes, as: Themes
require Logger
require Phoenix.PubSub
@impl true
def init(opts) do
children = [
add_pub_sub(opts),
{PartitionSupervisor, child_spec: DynamicSupervisor, name: ExESDB.EmitterPools},
{ExESDB.Store, opts},
{ExESDB.Cluster, opts},
{ExESDB.Streams, opts},
{ExESDB.Snapshots, opts},
{ExESDB.Subscriptions, opts},
{ExESDB.GatewaySupervisor, opts},
{ExESDB.LeaderSystem, opts}
]
:os.set_signal(:sigterm, :handle)
:os.set_signal(:sigquit, :handle)
spawn(fn -> handle_os_signal() end)
ret =
Supervisor.init(
children,
strategy: :one_for_one
)
IO.puts("#{Themes.system(self())} is UP")
ret
end
defp add_pub_sub(opts) do
pub_sub = Keyword.get(opts, :pub_sub)
case pub_sub do
nil ->
add_pub_sub([pub_sub: :native] ++ opts)
:native ->
{ExESDB.PubSub, opts}
_ ->
{Phoenix.PubSub, name: pub_sub}
end
end
defp handle_os_signal do
receive do
{:signal, :sigterm} ->
Logger.warning("SIGTERM received. Stopping ExESDB")
stop(:sigterm)
{:signal, :sigquit} ->
Logger.warning("SIGQUIT received. Stopping ExESDB")
stop(:sigquit)
msg ->
IO.puts("Unknown signal: #{inspect(msg)}")
Logger.warning("Received unknown signal: #{inspect(msg)}")
end
handle_os_signal()
end
def stop(_reason \\ :normal) do
Process.sleep(2_000)
Application.stop(:ex_esdb)
end
def start_link(opts),
do:
Supervisor.start_link(
__MODULE__,
opts,
name: __MODULE__
)
def start(opts) do
case start_link(opts) do
{:ok, pid} -> pid
{:error, {:already_started, pid}} -> pid
{:error, reason} -> raise "failed to start eventstores supervisor: #{inspect(reason)}"
end
end
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :supervisor
}
end
defp store(opts), do: Keyword.get(opts, :store, Options.store_id())
end