Packages

macula

0.10.2
7.1.0 7.0.0 6.0.0 5.2.2 5.2.1 5.2.0 5.1.0 5.0.0 4.8.0 4.7.1 4.7.0 4.6.0 4.5.0 4.4.10 4.4.9 4.4.8 4.4.7 4.4.6 4.4.5 4.4.4 4.4.3 4.4.2 4.4.1 4.4.0 4.3.1 4.3.0 4.2.9 4.2.8 4.2.7 4.2.6 4.2.5 4.2.4 4.2.3 4.2.2 4.2.1 4.2.0 4.1.1 4.1.0 4.0.0 3.16.0 3.15.3 3.15.2 3.15.1 3.14.0 3.13.0 3.12.1 3.12.0 3.11.1 3.11.0 3.10.3 3.10.2 3.10.1 3.9.0 3.8.0 3.7.0 3.5.0 3.4.0 3.3.0 3.2.0 3.1.0 3.0.0 2.1.1 2.1.0 2.0.0 1.5.2 1.5.1 1.4.30 1.4.29 1.4.28 1.4.27 1.4.26 1.4.25 1.4.24 1.4.23 1.4.22 1.4.21 1.4.20 1.4.19 1.4.18 1.4.17 1.4.16 1.4.15 1.4.14 1.4.13 1.4.11 1.4.10 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.1 1.3.0 1.2.0 1.1.0 1.0.10 1.0.9 1.0.8 1.0.7 1.0.6 1.0.5 1.0.4 1.0.3 1.0.2 1.0.1 1.0.0 0.48.6 0.48.5 0.48.4 0.48.3 0.48.2 0.48.1 0.48.0 0.47.1 0.47.0 0.46.3 0.46.1 0.46.0 0.45.3 0.45.2 0.45.1 0.45.0 0.44.2 0.44.1 0.44.0 0.43.3 0.43.2 0.43.1 0.43.0 0.42.9 0.42.8 0.42.7 0.42.6 0.42.5 0.42.4 0.42.3 0.42.2 0.42.1 0.42.0 0.41.1 0.41.0 0.40.1 0.40.0 0.39.9 0.39.8 0.39.7 0.39.6 0.39.5 0.39.4 0.39.3 0.39.2 0.39.1 0.39.0 0.38.8 0.38.7 0.38.6 0.38.5 0.38.4 0.38.3 0.38.2 0.38.1 0.38.0 0.37.7 0.37.6 0.37.5 0.37.4 0.37.3 0.37.2 0.37.1 0.37.0 0.36.6 0.36.5 0.36.4 0.36.3 0.36.2 0.36.1 0.36.0 0.35.4 0.35.3 0.35.2 0.35.1 0.35.0 0.34.1 0.34.0 0.33.1 0.33.0 0.32.5 0.32.4 0.32.3 0.32.2 0.32.1 0.32.0 0.31.9 0.31.8 0.31.7 0.31.6 0.31.5 0.31.4 0.31.3 0.31.2 0.31.1 0.31.0 0.30.10 0.30.9 0.30.8 0.30.7 0.30.6 0.30.5 0.30.4 0.30.3 0.30.2 0.30.1 0.30.0 0.29.0 0.28.3 0.28.2 0.28.1 0.28.0 0.27.1 0.27.0 0.26.1 0.26.0 0.25.6 0.25.5 0.25.4 0.25.3 0.25.2 0.25.1 0.25.0 0.24.6 0.24.5 0.24.4 0.24.3 0.24.2 0.24.1 0.24.0 0.23.3 0.23.2 0.23.1 0.23.0 0.22.12 0.22.11 0.22.10 0.22.9 0.22.8 0.22.7 0.22.6 0.22.5 0.22.4 0.22.3 0.22.2 0.22.1 0.22.0 0.21.7 0.21.6 0.21.5 0.21.4 0.21.2 0.21.1 0.21.0 0.20.25 0.20.24 0.20.23 0.20.22 0.20.21 0.20.20 0.20.19 0.20.18 0.20.17 0.20.16 0.20.15 0.20.14 0.20.13 0.20.12 0.20.11 0.20.10 0.20.9 0.20.8 0.20.7 0.20.6 0.20.5 0.20.3 0.20.2 0.20.1 0.20.0 0.19.2 0.19.1 0.19.0 0.18.1 0.18.0 0.17.4 0.17.3 0.17.2 0.17.1 0.17.0 0.16.6 0.16.5 0.16.4 0.16.3 0.16.2 0.16.1 0.16.0 0.15.1 0.15.0 0.14.3 0.14.2 0.14.1 0.14.0 0.12.6 0.12.5 0.12.3 0.11.3 0.10.2 0.10.1 0.10.0 0.9.2 0.9.1 0.9.0 0.8.25 0.8.24 0.8.23 0.8.22 0.8.21 0.8.20 0.8.19 0.8.18 0.8.17 0.8.16 0.8.15 0.8.14 0.8.13 0.8.12 0.8.11 0.8.10 0.8.9 0.8.8 0.8.7 0.8.6 0.8.5 0.8.4 0.8.3 0.8.2 0.8.1 0.8.0 0.7.30 0.7.29 0.7.28 0.7.27 0.7.26 0.7.25 0.7.24 0.7.23 0.7.22 0.7.21 0.7.20 0.7.19 0.7.18 0.7.17 0.7.16 0.7.15 0.7.14 0.7.13 0.7.12 0.7.11 0.7.10 0.7.9 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.7 0.6.6 0.6.5 0.6.4 0.6.3 0.6.2 0.6.1 0.6.0 0.5.0 0.4.4 0.4.3 0.4.2 0.4.1 0.4.0 0.3.4 0.3.3 0.3.2 0.3.1

Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise

Current section

Files

Jump to
macula src macula_platform_system macula_leader_election.erl
Raw

src/macula_platform_system/macula_leader_election.erl

%%%-------------------------------------------------------------------
%% @doc Macula Leader Election using Raft Consensus (via ra library).
%%
%% This module provides distributed leader election for platform services.
%% It uses the ra library (Raft implementation) to achieve consensus across
%% the mesh network.
%%
%% Use Cases:
%% - Single coordinator election for matchmaking
%% - Primary node selection for stateful services
%% - Distributed locking primitives
%%
%% Architecture:
%% - Uses ra (Raft) for consensus
%% - Discovers cluster members via Macula DHT
%% - Provides callbacks for leadership changes
%%
%% @end
%%%-------------------------------------------------------------------
-module(macula_leader_election).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([start_link/1]).
-export([get_leader/0, is_leader/0, get_members/0]).
-export([register_callback/2, unregister_callback/1]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
node_id :: binary(),
realm :: binary(),
cluster_name :: atom(),
server_id :: atom(),
leader :: atom() | undefined,
members :: [atom()],
callbacks :: #{atom() => fun()}
}).
-define(SERVER, ?MODULE).
-define(CLUSTER_NAME, macula_platform_cluster).
%%====================================================================
%% API functions
%%====================================================================
start_link(Config) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Config, []).
%% @doc Get the current leader node
get_leader() ->
gen_server:call(?SERVER, get_leader).
%% @doc Check if this node is the leader
is_leader() ->
gen_server:call(?SERVER, is_leader).
%% @doc Get all cluster members
get_members() ->
gen_server:call(?SERVER, get_members).
%% @doc Register a callback for leadership changes
%% Callback fun takes (IsLeader :: boolean())
register_callback(CallbackId, Fun) when is_function(Fun, 1) ->
gen_server:call(?SERVER, {register_callback, CallbackId, Fun}).
%% @doc Unregister a leadership callback
unregister_callback(CallbackId) ->
gen_server:call(?SERVER, {unregister_callback, CallbackId}).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init(Config) ->
NodeId = maps:get(node_id, Config),
Realm = maps:get(realm, Config),
?LOG_INFO("Initializing for node ~s in realm ~s",
[NodeId, Realm]),
%% Create unique server ID for this node
ServerIdBin = <<Realm/binary, "_", NodeId/binary>>,
ServerId = binary_to_atom(ServerIdBin, utf8),
State = #state{
node_id = NodeId,
realm = Realm,
cluster_name = ?CLUSTER_NAME,
server_id = ServerId,
leader = undefined,
members = [],
callbacks = #{}
},
%% Schedule cluster initialization after a delay
%% (allows mesh to stabilize)
erlang:send_after(5000, self(), init_cluster),
?LOG_INFO("Waiting for mesh to stabilize before initializing Raft cluster"),
{ok, State}.
handle_call(get_leader, _From, State) ->
{reply, {ok, State#state.leader}, State};
handle_call(is_leader, _From, State) ->
IsLeader = (State#state.server_id =:= State#state.leader),
{reply, IsLeader, State};
handle_call(get_members, _From, State) ->
{reply, {ok, State#state.members}, State};
handle_call({register_callback, CallbackId, Fun}, _From, State) ->
Callbacks = maps:put(CallbackId, Fun, State#state.callbacks),
%% Immediately notify new callback of current leadership status (Khepri pattern)
%% This ensures callbacks registered after election completes still get notified
case State#state.leader of
undefined ->
ok; % No leader yet, callback will fire when leader is elected
Leader ->
IsLeader = (State#state.server_id =:= Leader),
try
Fun(IsLeader)
catch
Class:Reason:Stacktrace ->
?LOG_ERROR("Callback ~p failed on registration: ~p:~p~n~p",
[CallbackId, Class, Reason, Stacktrace])
end
end,
{reply, ok, State#state{callbacks = Callbacks}};
handle_call({unregister_callback, CallbackId}, _From, State) ->
Callbacks = maps:remove(CallbackId, State#state.callbacks),
{reply, ok, State#state{callbacks = Callbacks}}.
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(init_cluster, State) ->
?LOG_INFO("Initializing Raft cluster"),
%% For initial implementation, create a single-node cluster
%% TODO: Discover peers via DHT and add them to cluster
ClusterName = State#state.cluster_name,
ServerId = State#state.server_id,
%% Start ra system if not already started
case ra:start() of
ok ->
?LOG_INFO("Ra system started");
{error, {already_started, _}} ->
?LOG_INFO("Ra system already running");
{error, Reason} ->
?LOG_ERROR("Failed to start ra: ~p", [Reason])
end,
%% Define the ra server configuration
%% Using our custom state machine for leader election
Machine = {module, macula_leader_machine, #{}},
%% Generate proper UID using ra:new_uid (Khepri pattern)
UId = ra:new_uid(atom_to_binary(ClusterName, utf8)),
ServerConfig = #{
id => ServerId,
uid => UId,
cluster_name => ClusterName,
log_init_args => #{uid => UId},
initial_members => [ServerId], % Single node for now
machine => Machine
},
%% Start the ra server
case ra:start_server(ServerConfig) of
ok ->
?LOG_INFO("Ra server started: ~p", [ServerId]),
%% Trigger election
case ra:trigger_election(ServerId) of
ok ->
?LOG_INFO("Election triggered");
{error, ElectReason} ->
?LOG_WARNING("Election trigger failed: ~p", [ElectReason])
end,
%% Start checking for leader immediately (Khepri pattern: aggressive polling initially)
erlang:send_after(500, self(), check_leader),
{noreply, State#state{members = [ServerId]}};
{error, {already_started, _}} ->
?LOG_INFO("Ra server already started"),
%% Start checking immediately
erlang:send_after(500, self(), check_leader),
{noreply, State#state{members = [ServerId]}};
{error, StartReason} ->
?LOG_ERROR("Failed to start ra server: ~p", [StartReason]),
%% Retry after delay
erlang:send_after(5000, self(), init_cluster),
{noreply, State}
end;
handle_info(check_leader, State) ->
ServerId = State#state.server_id,
%% Query ra for current leader with timeout (Khepri pattern)
NewLeader = case ra:members(ServerId, 2000) of
{ok, Members, Leader} ->
?LOG_DEBUG("Cluster state - Members: ~p, Leader: ~p",
[Members, Leader]),
OldLeader = State#state.leader,
%% Check if leadership changed
case {OldLeader, Leader} of
{Leader, Leader} ->
ok; % No change
{_, undefined} ->
?LOG_DEBUG("No leader elected yet, retrying...");
{undefined, Leader} when Leader =/= undefined ->
?LOG_INFO("Leader elected: ~p", [Leader]),
IsLeader = (ServerId =:= Leader),
notify_callbacks(IsLeader, State#state.callbacks);
{OldLeader, Leader} when OldLeader =/= Leader ->
?LOG_INFO("Leadership changed: ~p -> ~p",
[OldLeader, Leader]),
IsLeader = (ServerId =:= Leader),
notify_callbacks(IsLeader, State#state.callbacks)
end,
Leader;
{timeout, _} ->
?LOG_WARNING("Timeout querying ra members, will retry"),
State#state.leader;
{error, Reason} ->
?LOG_ERROR("Failed to query ra members: ~p, will retry", [Reason]),
State#state.leader
end,
%% Schedule next check (shorter interval if no leader yet - Khepri pattern)
NextCheckInterval = case NewLeader of
undefined -> 1000; % Check more frequently when waiting for leader
_ -> 5000 % Normal interval when leader is established
end,
erlang:send_after(NextCheckInterval, self(), check_leader),
{noreply, State#state{leader = NewLeader}};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%====================================================================
%% Internal functions
%%====================================================================
%% @private
notify_callbacks(IsLeader, Callbacks) ->
maps:foreach(
fun(CallbackId, Fun) ->
try
Fun(IsLeader)
catch
Class:Reason:Stacktrace ->
?LOG_ERROR("Callback ~p failed: ~p:~p~n~p",
[CallbackId, Class, Reason, Stacktrace])
end
end,
Callbacks
).