Packages
ex_esdb
0.4.0
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_system.ex
defmodule ExESDB.ClusterSystem do
@moduledoc """
Supervisor for cluster coordination components.
This supervisor manages cluster-specific coordination components:
- ClusterCoordinator: Handles coordination logic and split-brain prevention
- NodeMonitor: Monitors node health and handles failures
Note: KhepriCluster is managed at the System level since it's mode-aware.
"""
use Supervisor
require Logger
alias ExESDB.StoreNaming, as: StoreNaming
alias ExESDB.Themes, as: Themes
defp store_id(opts) do
Keyword.get(opts, :store_id, :ex_esdb)
end
@impl true
def init(opts) do
Logger.info("[CLUSTER_SYSTEM] Initializing ClusterSystem supervisor with opts: #{inspect(opts)}")
Logger.info("[CLUSTER_SYSTEM] ClusterSystem PID: #{inspect(self())}, Node: #{inspect(node())}")
# Extract store_id for logging
store_id = store_id(opts)
Logger.info("[CLUSTER_SYSTEM] Store ID: #{inspect(store_id)}")
Logger.info("[CLUSTER_SYSTEM] Starting cluster coordination components...")
children = [
# ClusterCoordinator handles coordination logic and split-brain prevention
{ExESDB.StoreCoordinator, opts},
# NodeMonitor provides fast failure detection
{ExESDB.NodeMonitor, node_monitor_config(opts)}
]
Logger.info("[CLUSTER_SYSTEM] Configured #{length(children)} cluster subsystems:")
children |> Enum.with_index(1) |> Enum.each(fn {{module, _opts}, index} ->
Logger.info("[CLUSTER_SYSTEM] #{index}. #{inspect(module)} - #{get_component_description(module)}")
end)
Logger.info("[CLUSTER_SYSTEM] Using :one_for_one strategy - independent cluster components")
IO.puts("#{Themes.cluster_system(self(), "is UP")}")
result = Supervisor.init(children, strategy: :one_for_one)
Logger.info("[CLUSTER_SYSTEM] ClusterSystem supervisor initialization complete")
result
end
defp get_component_description(ExESDB.StoreCoordinator), do: "Cluster coordination and split-brain prevention"
defp get_component_description(ExESDB.NodeMonitor), do: "Node health monitoring and failure detection"
defp get_component_description(_), do: "Cluster component"
def start_link(opts) do
Supervisor.start_link(
__MODULE__,
opts,
name: StoreNaming.genserver_name(__MODULE__, store_id(opts))
)
end
def child_spec(opts) do
%{
id: StoreNaming.child_spec_id(__MODULE__, store_id(opts)),
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: :infinity,
type: :supervisor
}
end
# Helper function to configure NodeMonitor options
defp node_monitor_config(opts) do
store_id = store_id(opts)
# More lenient configuration for node monitoring to prevent cascading failures
[
store_id: store_id,
# 5 seconds (less frequent probing)
probe_interval: 5_000,
# 6 consecutive failures (more tolerance)
failure_threshold: 6,
# 3 second timeout per probe (more time)
probe_timeout: 3_000
]
end
end