Packages
ex_esdb
0.0.13-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/cluster.ex
defmodule ExESDB.Cluster do
@moduledoc false
use GenServer
require Logger
alias ExESDB.LeaderWorker, as: LeaderWorker
alias ExESDB.Options, as: Opts
alias ExESDB.Themes, as: Themes
defp ping?(node) do
case :net_adm.ping(node) do
:pong -> true
_ -> false
end
end
def leader?(store) do
{_, leader_node} =
:ra_leaderboard.lookup_leader(store)
node() == leader_node
end
defp get_medal(leader, member),
do: if(member == leader, do: "🏆", else: "🥈")
defp join(store) do
Opts.seed_nodes()
|> Enum.map(fn seed ->
if ping?(seed) do
Logger.debug("#{Themes.cluster(node())} => Joining: #{inspect(seed)}")
store
|> :khepri_cluster.join(seed)
end
end)
end
defp leave(store) do
case store |> :khepri_cluster.reset() do
:ok ->
IO.puts("#{Themes.cluster(node())} => Left cluster")
:ok
{:error, reason} ->
Logger.error(
"#{Themes.cluster(node())} => Failed to leave cluster. reason: #{inspect(reason)}"
)
{:error, reason}
end
end
defp members(store),
do:
store
|> :khepri_cluster.members()
@impl true
def handle_info(:join, state) do
state[:store_id]
|> join()
{:noreply, state}
end
@impl true
def handle_info(:members, state) do
IO.puts("\nMEMBERS")
leader = Keyword.get(state, :current_leader)
store = state[:store_id]
case store
|> members() do
{:error, reason} ->
IO.puts("⚠️⚠️ Failed to get store members. reason: #{inspect(reason)} ⚠️⚠️")
{:ok, members} ->
members
|> Enum.each(fn {_store, member} ->
medal = get_medal(leader, member)
IO.puts("#{medal} #{inspect(member)}")
end)
end
Process.send_after(self(), :members, 5 * state[:timeout])
{:noreply, state}
end
@impl true
def handle_info(:check_leader, state) do
timeout = state[:timeout]
current_leader =
state
|> Keyword.get(:current_leader)
store =
state
|> Keyword.get(:store_id)
new_state =
case :ra_leaderboard.lookup_leader(store) do
{_, leader_node} ->
if node() == leader_node && current_leader != leader_node do
IO.puts("⚠️⚠️ FOLLOW THE LEADER! ⚠️⚠️")
store
|> LeaderWorker.activate()
end
state
|> Keyword.put(:current_leader, leader_node)
:undefined ->
IO.puts("⚠️⚠️ No leader found. ⚠️⚠️")
state
end
Process.send_after(self(), :check_leader, timeout)
{:noreply, new_state}
end
@impl true
def handle_info({:DOWN, _ref, :process, pid, reason}, state) do
state[:store_id]
|> leave()
IO.puts("🔻🔻 #{Themes.cluster(pid)} going down with reason: #{inspect(reason)} 🔻🔻")
{:noreply, state}
end
@impl true
def handle_info({:EXIT, pid, reason}, state) do
IO.puts("#{Themes.cluster(pid)} exited with reason: #{inspect(reason)}")
state[:store_id]
|> leave()
{:noreply, state}
end
@impl true
def handle_info(_, state) do
{:noreply, state}
end
############# PLUMBING #############
@impl true
def terminate(reason, state) do
Logger.warning("#{Themes.cluster(self())} terminating with reason: #{inspect(reason)}")
state[:store_id]
|> leave()
:ok
end
@impl true
def init(config) do
timeout = config[:timeout] || 1000
state = Keyword.put(config, :timeout, timeout)
IO.puts("#{Themes.cluster(self())} is UP")
Process.flag(:trap_exit, true)
Process.send_after(self(), :join, timeout)
Process.send_after(self(), :members, 10 * timeout)
Process.send_after(self(), :check_leader, timeout)
{:ok, state}
end
def start_link(opts),
do:
GenServer.start_link(
__MODULE__,
opts,
name: __MODULE__
)
def child_spec(opts),
do: %{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: 10_000,
type: :worker
}
end