Packages

macula

0.39.6
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_dist_system macula_cluster_static.erl
Raw

src/macula_dist_system/macula_cluster_static.erl

%%%-------------------------------------------------------------------
%%% @doc Macula Static Cluster Strategy.
%%%
%%% A simple cluster formation strategy that connects to a predefined
%%% list of nodes. Equivalent to libcluster's Cluster.Strategy.Epmd.
%%%
%%% == Configuration ==
%%%
%%% Start with a list of nodes to connect to:
%%%
%%% ```
%%% {ok, _Pid} = macula_cluster_static:start_link(#{
%%% nodes => ['node1@host1', 'node2@host2', 'node3@host3'],
%%% reconnect_interval => 5000 %% ms, default 5000
%%% }).
%%% '''
%%%
%%% Or from environment variables:
%%%
%%% ```
%%% %% CLUSTER_NODES=node1@host1,node2@host2,node3@host3
%%% {ok, _Pid} = macula_cluster_static:start_link(#{}).
%%% '''
%%%
%%% == Behavior ==
%%%
%%% - Attempts to connect to all configured nodes on startup
%%% - Monitors connected nodes for disconnect events
%%% - Automatically reconnects to disconnected nodes
%%% - Ignores self-connection attempts
%%% - Logs connection/disconnection events
%%%
%%% == Callbacks ==
%%%
%%% Register a callback module to receive cluster events:
%%%
%%% ```
%%% {ok, _Pid} = macula_cluster_static:start_link(#{
%%% nodes => [...],
%%% callback => self() %% PID or {Module, Function}
%%% }).
%%% %% Receives: {macula_cluster, nodeup, Node}
%%% %% Receives: {macula_cluster, nodedown, Node}
%%% '''
%%%
%%% @copyright 2026 Macula.io Apache-2.0
%%% @end
%%%-------------------------------------------------------------------
-module(macula_cluster_static).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/0,
start_link/1,
stop/0,
stop/1,
get_nodes/0,
get_nodes/1,
get_connected/0,
get_connected/1,
add_node/1,
add_node/2,
remove_node/1,
remove_node/2
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-define(SERVER, ?MODULE).
-define(DEFAULT_RECONNECT_INTERVAL, 5000).
-record(state, {
%% Configured nodes to connect to
nodes :: [atom()],
%% Currently connected nodes
connected :: sets:set(atom()),
%% Reconnect interval in milliseconds
reconnect_interval :: pos_integer(),
%% Reconnect timer reference
reconnect_timer :: reference() | undefined,
%% Callback for cluster events (pid or {M,F})
callback :: pid() | {module(), atom()} | undefined
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the static cluster strategy with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start the static cluster strategy with options.
%%
%% Options:
%% - nodes: List of node atoms to connect to
%% - reconnect_interval: Milliseconds between reconnect attempts (default 5000)
%% - callback: PID or {Module, Function} to receive cluster events
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Stop the static cluster strategy.
-spec stop() -> ok.
stop() ->
stop(?SERVER).
%% @doc Stop a named static cluster strategy.
-spec stop(atom() | pid()) -> ok.
stop(NameOrPid) ->
gen_server:stop(NameOrPid).
%% @doc Get the list of configured nodes.
-spec get_nodes() -> [atom()].
get_nodes() ->
get_nodes(?SERVER).
%% @doc Get the list of configured nodes from a named instance.
-spec get_nodes(atom() | pid()) -> [atom()].
get_nodes(NameOrPid) ->
gen_server:call(NameOrPid, get_nodes).
%% @doc Get the list of currently connected nodes.
-spec get_connected() -> [atom()].
get_connected() ->
get_connected(?SERVER).
%% @doc Get the list of connected nodes from a named instance.
-spec get_connected(atom() | pid()) -> [atom()].
get_connected(NameOrPid) ->
gen_server:call(NameOrPid, get_connected).
%% @doc Add a node to the cluster configuration.
-spec add_node(atom()) -> ok.
add_node(Node) ->
add_node(?SERVER, Node).
%% @doc Add a node to a named instance.
-spec add_node(atom() | pid(), atom()) -> ok.
add_node(NameOrPid, Node) ->
gen_server:call(NameOrPid, {add_node, Node}).
%% @doc Remove a node from the cluster configuration.
-spec remove_node(atom()) -> ok.
remove_node(Node) ->
remove_node(?SERVER, Node).
%% @doc Remove a node from a named instance.
-spec remove_node(atom() | pid(), atom()) -> ok.
remove_node(NameOrPid, Node) ->
gen_server:call(NameOrPid, {remove_node, Node}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init(Opts) ->
process_flag(trap_exit, true),
%% Get nodes from options or environment
Nodes = resolve_nodes(Opts),
ReconnectInterval = maps:get(reconnect_interval, Opts, ?DEFAULT_RECONNECT_INTERVAL),
Callback = maps:get(callback, Opts, undefined),
%% Subscribe to node events
ok = net_kernel:monitor_nodes(true),
State = #state{
nodes = Nodes,
connected = sets:new([{version, 2}]),
reconnect_interval = ReconnectInterval,
callback = Callback
},
?LOG_INFO("[macula_cluster_static] Started with ~p configured node(s)",
[length(Nodes)]),
%% Attempt initial connections
self() ! connect_all,
{ok, State}.
%% @private
handle_call(get_nodes, _From, State) ->
{reply, State#state.nodes, State};
handle_call(get_connected, _From, State) ->
{reply, sets:to_list(State#state.connected), State};
handle_call({add_node, Node}, _From, State) ->
NewNodes = lists:usort([Node | State#state.nodes]),
%% Trigger connection attempt
self() ! {connect_node, Node},
{reply, ok, State#state{nodes = NewNodes}};
handle_call({remove_node, Node}, _From, State) ->
NewNodes = lists:delete(Node, State#state.nodes),
%% Disconnect if connected
case sets:is_element(Node, State#state.connected) of
true ->
erlang:disconnect_node(Node);
false ->
ok
end,
NewConnected = sets:del_element(Node, State#state.connected),
{reply, ok, State#state{nodes = NewNodes, connected = NewConnected}};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
%% Initial connection attempt
handle_info(connect_all, State) ->
NewState = connect_to_all_nodes(State),
%% Schedule reconnection timer
Timer = schedule_reconnect(State#state.reconnect_interval),
{noreply, NewState#state{reconnect_timer = Timer}};
%% Connect to a specific node
handle_info({connect_node, Node}, State) ->
NewState = try_connect(Node, State),
{noreply, NewState};
%% Periodic reconnection
handle_info(reconnect, State) ->
NewState = reconnect_disconnected(State),
Timer = schedule_reconnect(State#state.reconnect_interval),
{noreply, NewState#state{reconnect_timer = Timer}};
%% Node joined the cluster
handle_info({nodeup, Node}, State) ->
case lists:member(Node, State#state.nodes) of
true ->
?LOG_INFO("[macula_cluster_static] Node ~p connected", [Node]),
NewConnected = sets:add_element(Node, State#state.connected),
notify_callback(State#state.callback, nodeup, Node),
{noreply, State#state{connected = NewConnected}};
false ->
%% Not a node we manage
{noreply, State}
end;
%% Node left the cluster
handle_info({nodedown, Node}, State) ->
case lists:member(Node, State#state.nodes) of
true ->
?LOG_WARNING("[macula_cluster_static] Node ~p disconnected", [Node]),
NewConnected = sets:del_element(Node, State#state.connected),
notify_callback(State#state.callback, nodedown, Node),
{noreply, State#state{connected = NewConnected}};
false ->
%% Not a node we manage
{noreply, State}
end;
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, State) ->
%% Cancel reconnect timer
cancel_timer(State#state.reconnect_timer),
%% Unsubscribe from node events
catch net_kernel:monitor_nodes(false),
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @private Resolve nodes from options or environment
-spec resolve_nodes(map()) -> [atom()].
resolve_nodes(Opts) ->
case maps:get(nodes, Opts, undefined) of
undefined ->
parse_env_nodes();
Nodes when is_list(Nodes) ->
[ensure_atom(N) || N <- Nodes]
end.
%% @private Parse nodes from CLUSTER_NODES environment variable
-spec parse_env_nodes() -> [atom()].
parse_env_nodes() ->
case os:getenv("CLUSTER_NODES") of
false ->
[];
"" ->
[];
NodesStr ->
Nodes = string:tokens(NodesStr, ","),
[ensure_atom(string:trim(N)) || N <- Nodes]
end.
%% @private Connect to all configured nodes
-spec connect_to_all_nodes(#state{}) -> #state{}.
connect_to_all_nodes(State) ->
lists:foldl(
fun(Node, AccState) ->
try_connect(Node, AccState)
end,
State,
State#state.nodes
).
%% @private Reconnect to disconnected nodes
-spec reconnect_disconnected(#state{}) -> #state{}.
reconnect_disconnected(State) ->
Disconnected = [N || N <- State#state.nodes,
not sets:is_element(N, State#state.connected)],
lists:foldl(
fun(Node, AccState) ->
try_connect(Node, AccState)
end,
State,
Disconnected
).
%% @private Try to connect to a single node
-spec try_connect(atom(), #state{}) -> #state{}.
try_connect(Node, State) when Node =:= node() ->
%% Don't connect to ourselves
State;
try_connect(Node, State) ->
case sets:is_element(Node, State#state.connected) of
true ->
%% Already connected
State;
false ->
?LOG_DEBUG("[macula_cluster_static] Attempting connection to ~p", [Node]),
case net_kernel:connect_node(Node) of
true ->
?LOG_INFO("[macula_cluster_static] Connected to ~p", [Node]),
NewConnected = sets:add_element(Node, State#state.connected),
notify_callback(State#state.callback, nodeup, Node),
State#state{connected = NewConnected};
false ->
?LOG_DEBUG("[macula_cluster_static] Failed to connect to ~p", [Node]),
State;
ignored ->
%% net_kernel not running
?LOG_WARNING("[macula_cluster_static] net_kernel not running, "
"cannot connect to ~p", [Node]),
State
end
end.
%% @private Schedule reconnection timer
-spec schedule_reconnect(pos_integer()) -> reference().
schedule_reconnect(Interval) ->
erlang:send_after(Interval, self(), reconnect).
%% @private Cancel timer if defined
-spec cancel_timer(reference() | undefined) -> ok.
cancel_timer(undefined) -> ok;
cancel_timer(Timer) ->
erlang:cancel_timer(Timer),
ok.
%% @private Ensure value is an atom
-spec ensure_atom(atom() | string() | binary()) -> atom().
ensure_atom(Value) when is_atom(Value) -> Value;
ensure_atom(Value) when is_list(Value) -> list_to_atom(Value);
ensure_atom(Value) when is_binary(Value) -> binary_to_atom(Value, utf8).
%% @private Notify callback of cluster event
-spec notify_callback(pid() | {module(), atom()} | undefined,
nodeup | nodedown, atom()) -> ok.
notify_callback(undefined, _Event, _Node) ->
ok;
notify_callback(Pid, Event, Node) when is_pid(Pid) ->
Pid ! {macula_cluster, Event, Node},
ok;
notify_callback({Module, Function}, Event, Node) ->
_ = catch Module:Function(Event, Node),
ok.