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
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 ->
macula_routing_server:store(Pid, Key, Value)
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 ->
macula_routing_server:delete_local(Pid, Key, user_requested),
ok
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 ->
find_value_result(macula_routing_server:find_value(Pid, Key, #{}))
end.
%% @private
find_value_result({ok, Value}) -> {ok, Value};
find_value_result({error, not_found}) -> {error, not_found};
find_value_result(Other) -> Other.
%% @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) ->
Key = {p, l, {dht_prefix_subscription, Prefix}},
subscribe_impl(Pid, Prefix, Key);
subscribe(_, _) ->
ok.
%% @private
subscribe_impl(Pid, _Prefix, Key) when Pid =:= self() ->
subscribe_self(Key);
subscribe_impl(Pid, Prefix, _Key) ->
Pid ! {subscribe_dht_prefix, Prefix},
ok.
%% @private
subscribe_self(Key) ->
case gproc:where(Key) of
undefined ->
gproc:reg(Key),
ok;
_Pid ->
ok
end.
%% @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) ->
case whereis(gproc) of
undefined ->
ok;
_Pid ->
do_notify_prefix_subscribers(Key, Message)
end.
%% @private
do_notify_prefix_subscribers(Key, Message) ->
Matches = gproc:select({l, p}, [{{{'_', '_', {dht_prefix_subscription, '$1'}}, '_', '_'},
[], ['$1']}]),
lists:foreach(fun(Prefix) -> notify_if_prefix_matches(Key, Prefix, Message) end, Matches).
%% @private
notify_if_prefix_matches(Key, Prefix, Message) ->
case binary:match(Key, Prefix) of
{0, _} ->
gproc:send({p, l, {dht_prefix_subscription, Prefix}}, Message);
_ ->
ok
end.