Packages

macula

5.2.1
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_dist_discovery.erl
Raw

src/macula_dist_system/macula_dist_discovery.erl

%%%-------------------------------------------------------------------
%%% @doc EPMD Replacement using Macula DHT Discovery.
%%%
%%% This module provides decentralized node discovery for Erlang distribution,
%%% replacing the centralized EPMD (Erlang Port Mapper Daemon).
%%%
%%% == How it Works ==
%%%
%%% Instead of registering with a local EPMD daemon on port 4369,
%%% nodes announce themselves via Macula's DHT (Distributed Hash Table):
%%%
%%% 1. On startup, nodes call register_node/2 to announce themselves
%%% 2. Other nodes find peers via lookup_node/1 which queries the DHT
%%% 3. Subscribers get notified of node join/leave events
%%%
%%% == Discovery Mechanisms ==
%%%
%%% - **mDNS**: For local network discovery (no bootstrap required)
%%% - **DHT**: For internet-scale discovery (requires bootstrap nodes)
%%% - **Bootstrap**: Initial DHT seeds from known nodes
%%%
%%% == Usage ==
%%%
%%% Register this node:
%%% ok = macula_dist_discovery:register_node('4433@192.168.1.100', 4433).
%%%
%%% Look up a node:
%%% {ok, #{host := Host, port := Port}} = macula_dist_discovery:lookup_node('4433@192.168.1.100').
%%%
%%% Subscribe to node events:
%%% ok = macula_dist_discovery:subscribe(self()).
%%% %% Receive: {node_discovered, Node, IP, Port}
%%% %% Receive: {node_lost, Node}
%%%
%%% @copyright 2025 Macula.io Apache-2.0
%%% @end
%%%-------------------------------------------------------------------
-module(macula_dist_discovery).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/0,
start_link/1,
register_node/2,
unregister_node/1,
lookup_node/1,
lookup_node/2,
list_nodes/0,
subscribe/1,
unsubscribe/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2,
code_change/3
]).
-define(SERVER, ?MODULE).
-define(DHT_PREFIX, <<"_dist.node.">>).
-define(DEFAULT_TTL, 300). % 5 minutes
-define(REFRESH_INTERVAL, 60000). % 1 minute
-define(CLEANUP_INTERVAL, 120000). % 2 minutes
-record(state, {
%% Local node registration
local_node :: atom() | undefined,
local_port :: integer() | undefined,
local_info :: map() | undefined,
%% Known nodes cache
nodes :: #{atom() => map()},
%% Subscribers for node events
subscribers :: [pid()],
%% Discovery type
discovery_type :: mdns | dht | both,
%% Timers
refresh_timer :: reference() | undefined,
cleanup_timer :: reference() | undefined
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the discovery server with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start the discovery server with options.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Register this node in the distributed registry.
%% This replaces EPMD registration.
-spec register_node(atom(), integer()) -> ok | {error, term()}.
register_node(NodeName, Port) ->
gen_server:call(?SERVER, {register_node, NodeName, Port}).
%% @doc Unregister this node from the distributed registry.
-spec unregister_node(atom()) -> ok.
unregister_node(NodeName) ->
gen_server:call(?SERVER, {unregister_node, NodeName}).
%% @doc Look up a node by name.
%% This replaces EPMD lookup.
-spec lookup_node(atom()) -> {ok, map()} | {error, not_found}.
lookup_node(NodeName) ->
lookup_node(NodeName, 5000).
%% @doc Look up a node by name with timeout.
-spec lookup_node(atom(), timeout()) -> {ok, map()} | {error, not_found}.
lookup_node(NodeName, Timeout) ->
gen_server:call(?SERVER, {lookup_node, NodeName}, Timeout).
%% @doc List all known nodes.
-spec list_nodes() -> [atom()].
list_nodes() ->
gen_server:call(?SERVER, list_nodes).
%% @doc Subscribe to node discovery events.
%% Subscriber will receive:
%% {node_discovered, NodeName, IP, Port}
%% {node_lost, NodeName}
-spec subscribe(pid()) -> ok.
subscribe(Pid) ->
gen_server:call(?SERVER, {subscribe, Pid}).
%% @doc Unsubscribe from node discovery events.
-spec unsubscribe(pid()) -> ok.
unsubscribe(Pid) ->
gen_server:call(?SERVER, {unsubscribe, Pid}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init(Opts) ->
process_flag(trap_exit, true),
DiscoveryType = maps:get(discovery_type, Opts, both),
State = #state{
nodes = #{},
subscribers = [],
discovery_type = DiscoveryType
},
%% Start refresh timer
RefreshTimer = erlang:send_after(?REFRESH_INTERVAL, self(), refresh_registration),
%% Start cleanup timer
CleanupTimer = erlang:send_after(?CLEANUP_INTERVAL, self(), cleanup_expired),
%% Subscribe to Macula DHT events if available
maybe_subscribe_to_dht(),
{ok, State#state{
refresh_timer = RefreshTimer,
cleanup_timer = CleanupTimer
}}.
%% @private
handle_call({register_node, NodeName, Port}, _From, State) ->
{ok, Host} = inet:gethostname(),
{ok, Addrs} = inet:getaddrs(Host, inet),
IP = hd(Addrs),
NodeInfo = #{
name => NodeName,
port => Port,
host => Host,
ip => IP,
protocol => 'macula-dist',
registered_at => erlang:system_time(second),
ttl => ?DEFAULT_TTL
},
%% Store in DHT
ok = store_in_dht(NodeName, NodeInfo),
%% Also announce via mDNS if enabled
maybe_announce_mdns(NodeName, Port, State),
{reply, ok, State#state{
local_node = NodeName,
local_port = Port,
local_info = NodeInfo
}};
handle_call({unregister_node, NodeName}, _From, State) ->
%% Remove from DHT
ok = remove_from_dht(NodeName),
%% Remove from mDNS if enabled
maybe_unannounce_mdns(NodeName, State),
NewState = case State#state.local_node of
NodeName ->
State#state{
local_node = undefined,
local_port = undefined,
local_info = undefined
};
_ ->
State
end,
{reply, ok, NewState};
handle_call({lookup_node, NodeName}, _From, State) ->
%% First check local cache
lookup_cached(maps:get(NodeName, State#state.nodes, undefined), NodeName, State);
handle_call(list_nodes, _From, State) ->
Nodes = maps:keys(State#state.nodes),
{reply, Nodes, State};
handle_call({subscribe, Pid}, _From, State) ->
%% Monitor the subscriber
erlang:monitor(process, Pid),
Subscribers = [Pid | State#state.subscribers],
{reply, ok, State#state{subscribers = Subscribers}};
handle_call({unsubscribe, Pid}, _From, State) ->
Subscribers = lists:delete(Pid, State#state.subscribers),
{reply, ok, State#state{subscribers = Subscribers}};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private Cache miss -> query DHT; cache hit -> validate
lookup_cached(undefined, NodeName, State) ->
%% Query DHT
lookup_miss_result(lookup_in_dht(NodeName), NodeName, State);
lookup_cached(NodeInfo, NodeName, State) ->
%% Check if cached entry is still valid
lookup_valid(is_entry_valid(NodeInfo), NodeInfo, NodeName, State).
lookup_miss_result({ok, NodeInfo}, NodeName, State) ->
%% Cache the result
NewNodes = maps:put(NodeName, NodeInfo, State#state.nodes),
{reply, {ok, NodeInfo}, State#state{nodes = NewNodes}};
lookup_miss_result({error, Reason}, _NodeName, State) ->
{reply, {error, Reason}, State}.
lookup_valid(true, NodeInfo, _NodeName, State) ->
{reply, {ok, NodeInfo}, State};
lookup_valid(false, _NodeInfo, NodeName, State) ->
%% Expired, refresh from DHT
lookup_refresh_result(lookup_in_dht(NodeName), NodeName, State).
lookup_refresh_result({ok, FreshInfo}, NodeName, State) ->
NewNodes = maps:put(NodeName, FreshInfo, State#state.nodes),
{reply, {ok, FreshInfo}, State#state{nodes = NewNodes}};
lookup_refresh_result({error, Reason}, NodeName, State) ->
NewNodes = maps:remove(NodeName, State#state.nodes),
{reply, {error, Reason}, State#state{nodes = NewNodes}}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info(refresh_registration, State) ->
%% Refresh our registration in DHT
NewState = case State#state.local_info of
undefined ->
State;
NodeInfo ->
ok = store_in_dht(State#state.local_node, NodeInfo),
State
end,
%% Reschedule
Timer = erlang:send_after(?REFRESH_INTERVAL, self(), refresh_registration),
{noreply, NewState#state{refresh_timer = Timer}};
handle_info(cleanup_expired, State) ->
%% Remove expired entries from cache
Now = erlang:system_time(second),
NewNodes = maps:filter(
fun(_NodeName, NodeInfo) ->
is_entry_valid(NodeInfo, Now)
end,
State#state.nodes
),
%% Notify subscribers about lost nodes
LostNodes = maps:keys(State#state.nodes) -- maps:keys(NewNodes),
lists:foreach(
fun(NodeName) ->
notify_subscribers({node_lost, NodeName}, State#state.subscribers)
end,
LostNodes
),
%% Reschedule
Timer = erlang:send_after(?CLEANUP_INTERVAL, self(), cleanup_expired),
{noreply, State#state{nodes = NewNodes, cleanup_timer = Timer}};
%% Handle DHT discovery events
handle_info({dht_node_discovered, NodeName, NodeInfo}, State) ->
%% Update cache
NewNodes = maps:put(NodeName, NodeInfo, State#state.nodes),
%% Notify subscribers
#{ip := IP, port := Port} = NodeInfo,
notify_subscribers({node_discovered, NodeName, IP, Port}, State#state.subscribers),
{noreply, State#state{nodes = NewNodes}};
%% Handle mDNS discovery events
handle_info({mdns_node_discovered, NodeName, IP, Port}, State) ->
NodeInfo = #{
name => NodeName,
port => Port,
ip => IP,
protocol => 'macula-dist',
registered_at => erlang:system_time(second),
ttl => ?DEFAULT_TTL,
source => mdns
},
NewNodes = maps:put(NodeName, NodeInfo, State#state.nodes),
notify_subscribers({node_discovered, NodeName, IP, Port}, State#state.subscribers),
{noreply, State#state{nodes = NewNodes}};
handle_info({'DOWN', _Ref, process, Pid, _Reason}, State) ->
%% Subscriber died, remove from list
Subscribers = lists:delete(Pid, State#state.subscribers),
{noreply, State#state{subscribers = Subscribers}};
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, State) ->
%% Unregister from DHT
case State#state.local_node of
undefined -> ok;
NodeName -> remove_from_dht(NodeName)
end,
%% Cancel timers
cancel_timer(State#state.refresh_timer),
cancel_timer(State#state.cleanup_timer),
ok.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal Functions - DHT Operations
%%%===================================================================
%% @private Store node info in DHT.
%% macula_routing_dht is an OPTIONAL runtime extension provided by
%% macula-relay; not part of the SDK. Calls go through erlang:apply/3
%% so dialyzer doesn't try to resolve the module statically.
store_in_dht(NodeName, NodeInfo) ->
Key = make_dht_key(NodeName),
case dht_available() of
false -> store_in_local_cache(NodeName, NodeInfo);
true ->
Ext = externalize_node_info(NodeInfo),
erlang:apply(macula_routing_dht, store, [Key, term_to_binary(Ext)]),
ok
end.
%% Atom-free wire form. The consumer decodes DHT values with
%% `binary_to_term(_, [safe])', which refuses terms containing atoms
%% unknown to the decoding node — a remote node's name atom usually
%% is. Ship `name'/`protocol' as binaries so the record stays
%% decodable everywhere.
externalize_node_info(#{name := Name, protocol := Proto} = NodeInfo) ->
NodeInfo#{name := atom_to_binary(Name, utf8),
protocol := atom_to_binary(Proto, utf8)}.
%% @private Remove node info from DHT
remove_from_dht(NodeName) ->
Key = make_dht_key(NodeName),
case dht_available() of
false -> remove_from_local_cache(NodeName);
true -> erlang:apply(macula_routing_dht, delete, [Key]), ok
end.
%% @private Look up node info in DHT
lookup_in_dht(NodeName) ->
Key = make_dht_key(NodeName),
case dht_available() of
false -> lookup_in_local_cache(NodeName);
true -> lookup_in_dht_result(erlang:apply(macula_routing_dht, find, [Key]), NodeName)
end.
%% @private Handle DHT lookup result.
%%
%% DHT values are attacker-influenceable (any peer that can write the
%% key controls the bytes), so the decode is `[safe]' — an unsafe
%% `binary_to_term' here allows permanent atom-table exhaustion and
%% crafted-term resource attacks. The try is required: `[safe]' has
%% no non-throwing variant and a hostile payload must degrade to a
%% lookup miss, not crash the discovery server.
lookup_in_dht_result({ok, BinInfo}, NodeName) ->
Decoded = try
{term, binary_to_term(BinInfo, [safe])}
catch
error:badarg -> undecodable
end,
validated_node_info(Decoded, NodeName);
lookup_in_dht_result({error, not_found}, NodeName) ->
lookup_in_local_cache(NodeName).
%% Reconstruct the internal node-info map from the atom-free wire
%% form, accepting only the expected shape. `name' is rebuilt from
%% the CALLER's NodeName — never atomized from DHT bytes (that would
%% reopen the atom-exhaustion vector via binary_to_atom).
validated_node_info({term, #{port := Port, host := Host,
registered_at := RegAt, ttl := Ttl} = Info},
NodeName)
when is_integer(Port), Port >= 0, Port =< 65535,
is_integer(RegAt), is_integer(Ttl), Ttl >= 0 ->
{ok, #{name => NodeName,
port => Port,
host => Host,
ip => maps:get(ip, Info, undefined),
protocol => 'macula-dist',
registered_at => RegAt,
ttl => Ttl}};
validated_node_info(_BadShape, NodeName) ->
?LOG_WARNING("[dist_discovery] Rejected malformed DHT node info for ~p",
[NodeName]),
lookup_in_local_cache(NodeName).
%% @private Make DHT key for node
make_dht_key(NodeName) when is_atom(NodeName) ->
<<?DHT_PREFIX/binary, (atom_to_binary(NodeName, utf8))/binary>>;
make_dht_key(NodeName) when is_list(NodeName) ->
<<?DHT_PREFIX/binary, (list_to_binary(NodeName))/binary>>.
%%%===================================================================
%%% Internal Functions - Local Cache (Fallback)
%%%===================================================================
%% @private Store in local ETS cache
store_in_local_cache(NodeName, NodeInfo) ->
ensure_local_cache(),
ets:insert(macula_dist_discovery_cache, {NodeName, NodeInfo}),
ok.
%% @private Remove from local ETS cache
remove_from_local_cache(NodeName) ->
ensure_local_cache(),
ets:delete(macula_dist_discovery_cache, NodeName),
ok.
%% @private Look up in local ETS cache
lookup_in_local_cache(NodeName) ->
ensure_local_cache(),
case ets:lookup(macula_dist_discovery_cache, NodeName) of
[{_, NodeInfo}] -> {ok, NodeInfo};
[] -> {error, not_found}
end.
%% @private Ensure local cache ETS table exists
ensure_local_cache() ->
case ets:info(macula_dist_discovery_cache) of
undefined ->
ets:new(macula_dist_discovery_cache, [
named_table,
public,
set,
{read_concurrency, true}
]);
_ ->
ok
end.
%%%===================================================================
%%% Internal Functions - mDNS
%%%===================================================================
%% @private Maybe announce via mDNS
maybe_announce_mdns(_NodeName, _Port, #state{discovery_type = dht}) ->
ok;
maybe_announce_mdns(NodeName, Port, _State) ->
case whereis(mdns_advertise_sup) of
undefined ->
ok;
_Pid ->
macula_dist_mdns_advertiser:register(NodeName, Port),
mdns_advertise_sup:start_child(macula_dist_mdns_advertiser),
ok
end.
%% @private Maybe unannounce via mDNS
maybe_unannounce_mdns(_NodeName, #state{discovery_type = dht}) ->
ok;
maybe_unannounce_mdns(_NodeName, _State) ->
case whereis(mdns_advertise_sup) of
undefined ->
ok;
_Pid ->
mdns_advertise:stop(macula_dist_mdns_advertiser),
macula_dist_mdns_advertiser:unregister(),
ok
end.
%%%===================================================================
%%% Internal Functions - DHT Subscription
%%%===================================================================
%% @private Subscribe to DHT events
maybe_subscribe_to_dht() ->
case dht_available() of
false -> ok;
true -> erlang:apply(macula_routing_dht, subscribe, [?DHT_PREFIX, self()])
end.
%% @private Check if macula_routing_dht module is available (loaded by macula-relay).
dht_available() ->
case code:ensure_loaded(macula_routing_dht) of
{module, _} -> whereis(macula_routing_dht) =/= undefined;
_ -> false
end.
%%%===================================================================
%%% Internal Functions - Utilities
%%%===================================================================
%% @private Check if entry is still valid
is_entry_valid(NodeInfo) ->
is_entry_valid(NodeInfo, erlang:system_time(second)).
is_entry_valid(NodeInfo, Now) ->
RegisteredAt = maps:get(registered_at, NodeInfo, 0),
TTL = maps:get(ttl, NodeInfo, ?DEFAULT_TTL),
(RegisteredAt + TTL * 2) > Now. % Grace period of 2x TTL
%% @private Notify all subscribers of an event
notify_subscribers(Event, Subscribers) ->
lists:foreach(
fun(Pid) ->
Pid ! Event
end,
Subscribers
).
%% @private Cancel timer if defined
cancel_timer(undefined) -> ok;
cancel_timer(Timer) -> erlang:cancel_timer(Timer).