Packages
ex_esdb
0.0.18
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/leader_worker.ex
defmodule ExESDB.LeaderWorker do
@moduledoc """
This module contains the leader's reponsibilities for the cluster.
"""
use GenServer
require Logger
alias ExESDB.Emitters
alias ExESDB.SubscriptionsReader, as: SubsR
alias ExESDB.Themes, as: Themes
############ API ############
def activate(store),
do:
GenServer.cast(
__MODULE__,
{:activate, store}
)
########## HANDLE_CAST ##########
@impl true
def handle_cast({:activate, store}, state) do
IO.puts("ππ Activating LEADER #{inspect(node())} ππ")
subscriptions =
store
|> SubsR.get_subscriptions()
case subscriptions
|> Enum.count() do
0 -> IO.puts("π¦π¦ No subscriptions found. π¦π¦")
num -> IO.puts("ππ #{num} subscriptions found. ππ")
end
subscriptions
|> Enum.each(fn {key, subscription} ->
IO.puts("ππ Starting Emitter for key #{inspect(key)} ππ")
store
|> Emitters.start_emitter(subscription)
end)
{:noreply, state}
end
@impl true
def handle_cast(msg, state) do
Logger.warning("Leader received unexpected CAST: #{inspect(msg)}")
{:noreply, state}
end
################ HANDLE_INFO ############
@impl true
def handle_info(msg, state) do
Logger.warning("Leader received unexpected INFO: #{inspect(msg)}")
{:noreply, state}
end
############# HANDLE_CALL ##########
@impl true
def handle_call(msg, _from, state) do
Logger.warning("Leader received unexpected CALL: #{inspect(msg)}")
{:reply, :ok, state}
end
############# PLUMBING #############
#
def start_link(opts),
do:
GenServer.start_link(
__MODULE__,
opts,
name: __MODULE__
)
@impl true
def terminate(reason, _state) do
Logger.warning("#{Themes.cluster(self())} terminating with reason: #{inspect(reason)}")
:ok
end
@impl true
def init(config) do
IO.puts("#{Themes.leader_worker(self())} is UP!")
Process.flag(:trap_exit, true)
{:ok, config}
end
def child_spec(opts),
do: %{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: 10_000,
type: :worker
}
end