Packages
macula
3.14.0
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
Current section
Files
src/macula_dist_system/macula_cluster_strategy.erl
%%%-------------------------------------------------------------------
%%% @doc Macula Cluster Strategy for libcluster.
%%%
%%% This module implements a cluster formation strategy using Macula's
%%% decentralized discovery (DHT/mDNS) instead of EPMD. It can be used
%%% with libcluster in Elixir or standalone in Erlang.
%%%
%%% Integration: Works with libcluster (Elixir) or standalone (Erlang).
%%%
%%% Discovery Modes:
%%% mdns - Local network discovery via mDNS (no bootstrap needed)
%%% dht - Internet-scale discovery via Macula DHT
%%% both - Try mDNS first, fall back to DHT
%%%
%%% Configuration: topology, realm, discovery_type
%%%
%%% @copyright 2025 Macula.io Apache-2.0
%%% @end
%%%-------------------------------------------------------------------
-module(macula_cluster_strategy).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/1,
start_link/2,
stop/1,
list_connected/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2,
code_change/3
]).
-define(POLL_INTERVAL, 5000). % 5 seconds
-define(CONNECT_TIMEOUT, 5000). % 5 seconds
-record(state, {
%% Topology name (for libcluster compatibility)
topology :: atom(),
%% Configuration
config :: map(),
%% Connected nodes
connected :: #{atom() => boolean()},
%% Discovery subscription reference
discovery_ref :: reference() | undefined,
%% Poll timer
poll_timer :: reference() | undefined,
%% Callback module (for libcluster integration)
callback_module :: module() | undefined
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the cluster strategy with options.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
Topology = maps:get(topology, Opts, macula_cluster),
start_link(Topology, Opts).
%% @doc Start the cluster strategy with name and options.
-spec start_link(atom(), map()) -> {ok, pid()} | {error, term()}.
start_link(Name, Opts) ->
gen_server:start_link({local, Name}, ?MODULE, Opts#{topology => Name}, []).
%% @doc Stop the cluster strategy.
-spec stop(atom() | pid()) -> ok.
stop(NameOrPid) ->
gen_server:stop(NameOrPid).
%% @doc List currently connected nodes.
-spec list_connected(atom() | pid()) -> [atom()].
list_connected(NameOrPid) ->
gen_server:call(NameOrPid, list_connected).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init(Opts) ->
process_flag(trap_exit, true),
Topology = maps:get(topology, Opts, macula_cluster),
Config = maps:get(config, Opts, #{}),
CallbackMod = maps:get(callback_module, Opts, undefined),
%% Subscribe to node discovery events
ok = macula_dist_discovery:subscribe(self()),
%% Subscribe to Erlang node monitoring
ok = net_kernel:monitor_nodes(true, [{node_type, all}]),
%% Start polling for existing nodes
PollTimer = erlang:send_after(?POLL_INTERVAL, self(), poll_nodes),
State = #state{
topology = Topology,
config = Config,
connected = #{},
poll_timer = PollTimer,
callback_module = CallbackMod
},
%% Log startup
?LOG_INFO(
"macula_cluster_strategy: started for topology ~p",
[Topology]
),
{ok, State}.
%% @private
handle_call(list_connected, _From, State) ->
Connected = maps:keys(maps:filter(fun(_, V) -> V end, State#state.connected)),
{reply, Connected, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
%% Handle node discovered via Macula discovery
handle_info({node_discovered, _NodeName, IP, Port}, State) ->
Node = make_node_name(Port, IP),
NewState = maybe_connect_node(Node, State),
{noreply, NewState};
%% Handle node lost via Macula discovery
handle_info({node_lost, NodeName}, State) ->
Node = ensure_atom(NodeName),
NewState = maybe_disconnect_node(Node, State),
{noreply, NewState};
%% Handle Erlang node up event
handle_info({nodeup, Node, _Info}, State) ->
?LOG_INFO(
"macula_cluster_strategy: node ~p joined cluster",
[Node]
),
Connected = maps:put(Node, true, State#state.connected),
notify_callback(State, {nodeup, Node}),
{noreply, State#state{connected = Connected}};
%% Handle Erlang node down event
handle_info({nodedown, Node, _Info}, State) ->
?LOG_INFO(
"macula_cluster_strategy: node ~p left cluster",
[Node]
),
Connected = maps:put(Node, false, State#state.connected),
notify_callback(State, {nodedown, Node}),
{noreply, State#state{connected = Connected}};
%% Periodic polling for nodes
handle_info(poll_nodes, State) ->
%% Query discovery for known nodes
NewState = poll_discovered_nodes(State),
%% Reschedule
Timer = erlang:send_after(?POLL_INTERVAL, self(), poll_nodes),
{noreply, NewState#state{poll_timer = Timer}};
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, State) ->
%% Unsubscribe from discovery
catch macula_dist_discovery:unsubscribe(self()),
%% Stop node monitoring
catch net_kernel:monitor_nodes(false),
%% Cancel timer
cancel_timer(State#state.poll_timer),
ok.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @private Try to connect to a discovered node
maybe_connect_node(Node, State) when Node =:= node() ->
%% Don't connect to ourselves
State;
maybe_connect_node(Node, State) ->
case maps:get(Node, State#state.connected, false) of
true ->
%% Already connected
State;
false ->
%% Try to connect
?LOG_INFO(
"macula_cluster_strategy: attempting to connect to ~p",
[Node]
),
case net_kernel:connect_node(Node) of
true ->
?LOG_INFO(
"macula_cluster_strategy: connected to ~p",
[Node]
),
Connected = maps:put(Node, true, State#state.connected),
notify_callback(State, {connected, Node}),
State#state{connected = Connected};
false ->
?LOG_WARNING(
"macula_cluster_strategy: failed to connect to ~p",
[Node]
),
State;
ignored ->
%% net_kernel not running
State
end
end.
%% @private Disconnect from a node that's no longer in discovery
maybe_disconnect_node(Node, State) ->
case maps:get(Node, State#state.connected, false) of
true ->
?LOG_INFO(
"macula_cluster_strategy: disconnecting from ~p (no longer in discovery)",
[Node]
),
erlang:disconnect_node(Node),
Connected = maps:put(Node, false, State#state.connected),
notify_callback(State, {disconnected, Node}),
State#state{connected = Connected};
false ->
State
end.
%% @private Poll discovered nodes
poll_discovered_nodes(State) ->
case macula_dist_discovery:list_nodes() of
Nodes when is_list(Nodes) ->
lists:foldl(
fun(NodeName, AccState) ->
case macula_dist_discovery:lookup_node(NodeName) of
{ok, #{ip := IP, port := Port}} ->
Node = make_node_name(Port, IP),
maybe_connect_node(Node, AccState);
{error, _} ->
AccState
end
end,
State,
Nodes
);
_ ->
State
end.
%% @private Make node name from port and IP
%% Format: port@ip (e.g., '4433@192.168.1.100')
make_node_name(Port, IP) when is_tuple(IP) ->
make_node_name(Port, inet:ntoa(IP));
make_node_name(Port, IP) when is_list(IP) ->
list_to_atom(integer_to_list(Port) ++ "@" ++ IP);
make_node_name(Port, IP) when is_binary(IP) ->
make_node_name(Port, binary_to_list(IP)).
%% @private Ensure value is an 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 module (for libcluster integration)
notify_callback(#state{callback_module = undefined}, _Event) ->
ok;
notify_callback(#state{callback_module = Mod, topology = Topology}, Event) ->
_ = catch Mod:handle_event(Topology, Event),
ok.
%% @private Cancel timer if defined
cancel_timer(undefined) -> ok;
cancel_timer(Timer) -> erlang:cancel_timer(Timer).