Packages
macula
0.16.5
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_routing_system/macula_routing_dht.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Core DHT algorithms for Kademlia routing.
%%% Implements iterative lookup, store, and find operations.
%%% Pure functions - no GenServer, designed to be called by macula_routing_server.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_routing_dht).
%% API - Core DHT algorithms
-export([
iterative_find_node/4,
store_value/6,
find_value/4,
update_closest/4,
select_alpha/3
]).
%% API - Simple key-value interface (used by macula_dist_discovery)
-export([
store/2,
delete/1,
find/1,
subscribe/2,
notify_store/2,
notify_delete/1
]).
%% Types
-type query_fn() :: fun((macula_routing_bucket:node_info(), binary()) ->
{ok, [macula_routing_bucket:node_info()]} |
{value, term()} |
{nodes, [macula_routing_bucket:node_info()]} |
{error, term()}).
-type store_fn() :: fun((macula_routing_bucket:node_info(), binary(), term()) -> ok | {error, term()}).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Iterative lookup to find k closest nodes to target.
%% Uses alpha concurrent queries (default: 3).
-spec iterative_find_node(
macula_routing_table:routing_table(),
binary(),
pos_integer(),
query_fn()
) -> {ok, [macula_routing_bucket:node_info()]}.
iterative_find_node(RoutingTable, Target, K, QueryFn) ->
%% 1. Get k closest from local routing table
InitialClosest = macula_routing_table:find_closest(RoutingTable, Target, K),
%% 2. Perform iterative lookup with alpha=3 concurrency
Alpha = 3,
FinalClosest = iterative_lookup(InitialClosest, Target, K, Alpha, [], QueryFn),
{ok, FinalClosest}.
%% @doc Store value at k closest nodes to key.
-spec store_value(
macula_routing_table:routing_table(),
binary(),
term(),
pos_integer(),
query_fn(),
store_fn()
) -> ok.
store_value(RoutingTable, Key, Value, K, QueryFn, StoreFn) ->
%% 1. Find k closest nodes to key
{ok, ClosestNodes} = iterative_find_node(RoutingTable, Key, K, QueryFn),
%% 2. Store value at each of the k closest
lists:foreach(
fun(Node) ->
StoreFn(Node, Key, Value)
end,
ClosestNodes
),
ok.
%% @doc Find value in DHT.
%% Returns {ok, Value} if found, {nodes, [NodeInfo]} if not found.
-spec find_value(
macula_routing_table:routing_table(),
binary(),
pos_integer(),
query_fn()
) -> {ok, term()} | {nodes, [macula_routing_bucket:node_info()]}.
find_value(RoutingTable, Key, K, QueryFn) ->
%% Get initial closest nodes
InitialClosest = macula_routing_table:find_closest(RoutingTable, Key, K),
%% Perform iterative lookup, but stop if value found
Alpha = 3,
iterative_find_value(InitialClosest, Key, K, Alpha, [], QueryFn).
%% @doc Update closest set with new nodes, maintaining k closest and removing duplicates.
-spec update_closest(
[macula_routing_bucket:node_info()],
[macula_routing_bucket:node_info()],
binary(),
pos_integer()
) -> [macula_routing_bucket:node_info()].
update_closest(CurrentClosest, NewNodes, Target, K) ->
Combined = CurrentClosest ++ NewNodes,
Deduplicated = deduplicate_nodes(Combined),
sort_by_distance_and_take(Deduplicated, Target, K).
sort_by_distance_and_take(Nodes, Target, K) ->
WithDistance = [{distance_to(Target, N), N} || N <- Nodes],
Sorted = lists:keysort(1, WithDistance),
[Node || {_Dist, Node} <- lists:sublist(Sorted, K)].
distance_to(Target, #{node_id := NodeId}) ->
macula_routing_nodeid:distance(Target, NodeId).
%% @doc Select up to alpha unqueried nodes from closest set.
-spec select_alpha(
[macula_routing_bucket:node_info()],
[binary()],
pos_integer()
) -> [macula_routing_bucket:node_info()].
select_alpha(Closest, Queried, Alpha) ->
QueriedSet = sets:from_list(Queried),
Unqueried = [N || #{node_id := Id} = N <- Closest, not sets:is_element(Id, QueriedSet)],
lists:sublist(Unqueried, Alpha).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Iterative lookup algorithm for FIND_NODE.
-spec iterative_lookup(
[macula_routing_bucket:node_info()],
binary(),
pos_integer(),
pos_integer(),
[binary()],
query_fn()
) -> [macula_routing_bucket:node_info()].
iterative_lookup(Closest, Target, K, Alpha, Queried, QueryFn) ->
ToQuery = select_alpha(Closest, Queried, Alpha),
do_iterative_lookup(ToQuery, Closest, Target, K, Alpha, Queried, QueryFn).
%% No more nodes to query - return current closest
do_iterative_lookup([], Closest, _Target, _K, _Alpha, _Queried, _QueryFn) ->
Closest;
%% Query nodes and continue if closer found
do_iterative_lookup(ToQuery, Closest, Target, K, Alpha, Queried, QueryFn) ->
{NewNodes, NewQueried} = query_nodes(ToQuery, Target, Queried, QueryFn),
UpdatedClosest = update_closest(Closest, NewNodes, Target, K),
continue_if_closer(Closest, UpdatedClosest, Target, K, Alpha, NewQueried, QueryFn).
continue_if_closer(OldClosest, NewClosest, Target, K, Alpha, Queried, QueryFn) ->
case found_closer_nodes(OldClosest, NewClosest, K) of
true -> iterative_lookup(NewClosest, Target, K, Alpha, Queried, QueryFn);
false -> NewClosest
end.
%% @doc Iterative lookup for FIND_VALUE (stops when value found).
-spec iterative_find_value(
[macula_routing_bucket:node_info()],
binary(),
pos_integer(),
pos_integer(),
[binary()],
query_fn()
) -> {ok, term()} | {nodes, [macula_routing_bucket:node_info()]}.
iterative_find_value(Closest, Key, K, Alpha, Queried, QueryFn) ->
ToQuery = select_alpha(Closest, Queried, Alpha),
do_iterative_find_value(ToQuery, Closest, Key, K, Alpha, Queried, QueryFn).
%% No more nodes to query
do_iterative_find_value([], Closest, _Key, _K, _Alpha, _Queried, _QueryFn) ->
{nodes, Closest};
%% Query nodes for value
do_iterative_find_value(ToQuery, Closest, Key, K, Alpha, Queried, QueryFn) ->
handle_value_query_result(
query_nodes_for_value(ToQuery, Key, Queried, QueryFn),
Closest, Key, K, Alpha, QueryFn
).
handle_value_query_result({value, Value}, _Closest, _Key, _K, _Alpha, _QueryFn) ->
{ok, Value};
handle_value_query_result({nodes, NewNodes, NewQueried}, Closest, Key, K, Alpha, QueryFn) ->
UpdatedClosest = update_closest(Closest, NewNodes, Key, K),
case found_closer_nodes(Closest, UpdatedClosest, K) of
true -> iterative_find_value(UpdatedClosest, Key, K, Alpha, NewQueried, QueryFn);
false -> {nodes, UpdatedClosest}
end.
%% @doc Query nodes and collect responses.
-spec query_nodes(
[macula_routing_bucket:node_info()],
binary(),
[binary()],
query_fn()
) -> {[macula_routing_bucket:node_info()], [binary()]}.
query_nodes(Nodes, Target, Queried, QueryFn) ->
lists:foldl(
fun(Node, {AccNodes, AccQueried}) ->
NodeId = maps:get(node_id, Node),
case QueryFn(Node, Target) of
{ok, ResponseNodes} ->
{AccNodes ++ ResponseNodes, [NodeId | AccQueried]};
{error, _Reason} ->
%% Query failed, just mark as queried
{AccNodes, [NodeId | AccQueried]}
end
end,
{[], Queried},
Nodes
).
%% @doc Query nodes for value (FIND_VALUE).
-spec query_nodes_for_value(
[macula_routing_bucket:node_info()],
binary(),
[binary()],
query_fn()
) -> {value, term()} | {nodes, [macula_routing_bucket:node_info()], [binary()]}.
query_nodes_for_value([], _Key, Queried, _QueryFn) ->
{nodes, [], Queried};
query_nodes_for_value([Node | Rest], Key, Queried, QueryFn) ->
NodeId = maps:get(node_id, Node),
case QueryFn(Node, Key) of
{value, Value} ->
%% Found the value!
{value, Value};
{nodes, ResponseNodes} ->
%% Continue querying
case query_nodes_for_value(Rest, Key, [NodeId | Queried], QueryFn) of
{value, Value} ->
{value, Value};
{nodes, AccNodes, AccQueried} ->
{nodes, ResponseNodes ++ AccNodes, AccQueried}
end;
{error, _Reason} ->
%% Query failed, try next node
query_nodes_for_value(Rest, Key, [NodeId | Queried], QueryFn)
end.
%% @doc Check if we found closer nodes (convergence check).
-spec found_closer_nodes(
[macula_routing_bucket:node_info()],
[macula_routing_bucket:node_info()],
pos_integer()
) -> boolean().
found_closer_nodes(OldClosest, NewClosest, K) ->
OldIds = extract_node_ids(lists:sublist(OldClosest, K)),
NewIds = extract_node_ids(lists:sublist(NewClosest, K)),
OldIds =/= NewIds.
extract_node_ids(Nodes) ->
lists:sort([Id || #{node_id := Id} <- Nodes]).
%% @doc Remove duplicate nodes (by node_id).
-spec deduplicate_nodes([macula_routing_bucket:node_info()]) -> [macula_routing_bucket:node_info()].
deduplicate_nodes(Nodes) ->
{_, Result} = lists:foldl(fun dedupe_node/2, {#{}, []}, Nodes),
lists:reverse(Result).
dedupe_node(#{node_id := Id} = Node, {Seen, Acc}) ->
case maps:is_key(Id, Seen) of
true -> {Seen, Acc};
false -> {Seen#{Id => true}, [Node | Acc]}
end.
%%%===================================================================
%%% Simple Key-Value Interface
%%% Used by macula_dist_discovery for QUIC distribution (deferred v1.1.0+)
%%%===================================================================
%% @doc Store a key-value pair in the DHT.
%% Delegates to macula_routing_server if running.
-spec store(binary(), binary()) -> ok | {error, term()}.
store(Key, Value) ->
case whereis(macula_routing_server) of
undefined ->
{error, not_started};
Pid ->
try
macula_routing_server:store(Pid, Key, Value)
catch
exit:{noproc, _} -> {error, not_started};
_:Reason -> {error, Reason}
end
end.
%% @doc Delete a key from the DHT.
%% Delegates to macula_routing_server if running.
-spec delete(binary()) -> ok | {error, term()}.
delete(Key) ->
case whereis(macula_routing_server) of
undefined ->
{error, not_started};
Pid ->
try
macula_routing_server:delete_local(Pid, Key, user_requested),
ok
catch
exit:{noproc, _} -> {error, not_started};
_:Reason -> {error, Reason}
end
end.
%% @doc Find a value in the DHT.
%% Delegates to macula_routing_server if running.
-spec find(binary()) -> {ok, binary()} | {error, not_found | term()}.
find(Key) ->
case whereis(macula_routing_server) of
undefined ->
{error, not_started};
Pid ->
try
case macula_routing_server:find_value(Pid, Key, #{}) of
{ok, Value} -> {ok, Value};
{error, not_found} -> {error, not_found};
Other -> Other
end
catch
exit:{noproc, _} -> {error, not_started};
_:Reason -> {error, Reason}
end
end.
%% @doc Subscribe to DHT events matching a key prefix.
%% Subscribes the given Pid to receive events when keys with the given prefix
%% are stored or deleted. Uses gproc property-based subscriptions.
%%
%% Events sent to the subscriber:
%% {dht_stored, Key, Value} - When a key is stored
%% {dht_deleted, Key} - When a key is deleted
%%
%% To unsubscribe, the subscriber process should call:
%% gproc:unreg({p, l, {dht_prefix_subscription, Prefix}})
-spec subscribe(binary(), pid()) -> ok.
subscribe(Prefix, Pid) when is_binary(Prefix), is_pid(Pid) ->
%% Register subscription via gproc property
%% The macula_routing_server will notify subscribers when keys change
Key = {p, l, {dht_prefix_subscription, Prefix}},
try
%% If the caller is subscribing itself, use gproc:reg
%% If subscribing another process, we need to send a message to that process
case Pid =:= self() of
true ->
gproc:reg(Key),
ok;
false ->
%% For remote subscription, send a message to the target pid
%% asking it to register. The target process must handle this.
Pid ! {subscribe_dht_prefix, Prefix},
ok
end
catch
error:badarg ->
%% Already registered - this is fine
ok;
_:_ ->
ok
end;
subscribe(_, _) ->
ok.
%% @doc Notify subscribers about a DHT store event.
%% Called by macula_routing_server when a key is stored.
-spec notify_store(binary(), term()) -> ok.
notify_store(Key, Value) ->
notify_prefix_subscribers(Key, {dht_stored, Key, Value}).
%% @doc Notify subscribers about a DHT delete event.
%% Called by macula_routing_server when a key is deleted.
-spec notify_delete(binary()) -> ok.
notify_delete(Key) ->
notify_prefix_subscribers(Key, {dht_deleted, Key}).
%% @doc Send notification to all subscribers whose prefix matches the key.
-spec notify_prefix_subscribers(binary(), term()) -> ok.
notify_prefix_subscribers(Key, Message) ->
%% Find all prefix subscriptions and check if they match
%% This iterates over all subscriptions - for production use with many
%% subscriptions, consider a more efficient data structure
try
Matches = gproc:select({l, p}, [{{{'_', '_', {dht_prefix_subscription, '$1'}}, '_', '_'},
[], ['$1']}]),
lists:foreach(fun(Prefix) ->
case binary:match(Key, Prefix) of
{0, _} ->
%% Key starts with this prefix - notify subscribers
gproc:send({p, l, {dht_prefix_subscription, Prefix}}, Message);
_ ->
ok
end
end, Matches)
catch
_:_ ->
%% gproc not available or no subscribers - ignore
ok
end.