Packages

macula

0.20.21
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).
%% 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
case maps:get(NodeName, State#state.nodes, undefined) of
undefined ->
%% Query DHT
case lookup_in_dht(NodeName) of
{ok, NodeInfo} ->
%% Cache the result
NewNodes = maps:put(NodeName, NodeInfo, State#state.nodes),
{reply, {ok, NodeInfo}, State#state{nodes = NewNodes}};
{error, Reason} ->
{reply, {error, Reason}, State}
end;
NodeInfo ->
%% Check if cached entry is still valid
case is_entry_valid(NodeInfo) of
true ->
{reply, {ok, NodeInfo}, State};
false ->
%% Expired, refresh from DHT
case lookup_in_dht(NodeName) of
{ok, FreshInfo} ->
NewNodes = maps:put(NodeName, FreshInfo, State#state.nodes),
{reply, {ok, FreshInfo}, State#state{nodes = NewNodes}};
{error, Reason} ->
NewNodes = maps:remove(NodeName, State#state.nodes),
{reply, {error, Reason}, State#state{nodes = NewNodes}}
end
end
end;
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
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
store_in_dht(NodeName, NodeInfo) ->
Key = make_dht_key(NodeName),
case whereis(macula_routing_dht) of
undefined ->
store_in_local_cache(NodeName, NodeInfo);
_Pid ->
macula_routing_dht:store(Key, term_to_binary(NodeInfo)),
ok
end.
%% @private Remove node info from DHT
remove_from_dht(NodeName) ->
Key = make_dht_key(NodeName),
case whereis(macula_routing_dht) of
undefined ->
remove_from_local_cache(NodeName);
_Pid ->
macula_routing_dht:delete(Key),
ok
end.
%% @private Look up node info in DHT
lookup_in_dht(NodeName) ->
Key = make_dht_key(NodeName),
case whereis(macula_routing_dht) of
undefined ->
lookup_in_local_cache(NodeName);
_Pid ->
lookup_in_dht_result(macula_routing_dht:find(Key), NodeName)
end.
%% @private Handle DHT lookup result
lookup_in_dht_result({ok, BinInfo}, _NodeName) ->
{ok, binary_to_term(BinInfo)};
lookup_in_dht_result({error, not_found}, 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 whereis(macula_routing_dht) of
undefined ->
ok;
_Pid ->
macula_routing_dht:subscribe(?DHT_PREFIX, self())
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).