Packages
ex_esdb
0.1.4
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.SubscriptionsWriter, as: SubsW
alias ExESDB.Themes, as: Themes
############ API ############
@doc """
Activates the LeaderWorker for the given store.
This function is called when this node becomes the cluster leader.
"""
def activate(store) do
# Save default subscriptions synchronously
case GenServer.call(__MODULE__, {:save_default_subscriptions, store}, 10_000) do
{:ok, _result} ->
# Now activate leadership responsibilities
GenServer.cast(__MODULE__, {:activate, store})
IO.puts(Themes.leader_worker(self(), "✅ LeaderWorker activated successfully"))
:ok
{:error, reason} ->
IO.puts(Themes.leader_worker(self(), "❌ Failed to activate LeaderWorker: #{inspect(reason)}"))
{:error, reason}
end
rescue
error ->
IO.puts(Themes.leader_worker(self(), "❌ LeaderWorker activation failed with error: #{inspect(error)}"))
{:error, error}
catch
:exit, reason ->
IO.puts(Themes.leader_worker(self(), "❌ LeaderWorker activation failed with exit: #{inspect(reason)}"))
{:error, reason}
end
########## HANDLE_CAST ##########
@impl true
def handle_cast({:activate, store}, state) do
IO.puts("\n#{Themes.leader_worker(self(), "🚀 ACTIVATING LEADERSHIP RESPONSIBILITIES")}")
IO.puts(" 🏆 Node: #{inspect(node())}")
IO.puts(" 📊 Store: #{inspect(store)}")
subscriptions =
store
|> SubsR.get_subscriptions()
subscription_count = Enum.count(subscriptions)
case subscription_count do
0 ->
IO.puts(" 📝 No active subscriptions to manage")
1 ->
IO.puts(" 📝 Managing 1 active subscription")
num ->
IO.puts(" 📝 Managing #{num} active subscriptions")
end
if subscription_count > 0 do
IO.puts("\n Starting emitters for active subscriptions:")
# Check if EmitterPools is available
case Process.whereis(ExESDB.EmitterPools) do
nil ->
IO.puts(" ⚠️ EmitterPools not available, skipping emitter startup")
IO.puts(" ℹ️ EmitterPools will be started when EmitterSystem is ready")
_pid ->
subscriptions
|> Enum.each(fn {key, subscription} ->
IO.puts(" ⚙️ Starting emitter for: #{inspect(key)}")
case Emitters.start_emitter_pool(store, subscription) do
{:ok, _pid} ->
IO.puts(" ✅ Successfully started emitter pool for #{inspect(key)}")
{:error, {:already_started, _pid}} ->
IO.puts(" ✅ Emitter pool already running for #{inspect(key)}")
{:error, reason} ->
IO.puts(" ❌ Failed to start emitter pool for #{inspect(key)}: #{inspect(reason)}")
end
end)
end
end
IO.puts("\n ✅ Leadership activation complete\n")
{: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({:save_default_subscriptions, store}, _from, state) do
try do
res =
store
|> SubsW.put_subscription_sync(:by_stream, "$all", "all-events")
{:reply, res, state}
rescue
error ->
Logger.warning("Failed to save default subscriptions: #{inspect(error)}")
{:reply, {:error, error}, state}
catch
:exit, reason ->
Logger.warning("Failed to save default subscriptions (exit): #{inspect(reason)}")
{:reply, {:error, reason}, state}
end
end
@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
# Set trap_exit early to handle crashes properly
Process.flag(:trap_exit, true)
# Log startup with process info
IO.puts("#{Themes.leader_worker(self(), "is UP!")}")
Logger.info("LeaderWorker started successfully with PID #{inspect(self())} and name #{inspect(__MODULE__)}")
# Verify we're properly registered
case Process.whereis(__MODULE__) do
pid when pid == self() ->
Logger.info("LeaderWorker registration confirmed: #{inspect(__MODULE__)} -> #{inspect(self())}")
other_pid ->
Logger.warning("LeaderWorker registration issue: expected #{inspect(self())}, got #{inspect(other_pid)}")
end
{:ok, config}
end
def child_spec(opts),
do: %{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: 10_000,
type: :worker
}
end