Packages
macula
0.10.2
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_server.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% GenServer managing Kademlia DHT routing table and operations.
%%% Integrates all routing components: table, DHT algorithms, protocol.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_routing_server).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/2,
add_node/2,
find_closest/3,
store_local/3,
store/3,
get_local/2,
get_all_keys/1,
find_value/3,
get_routing_table/1,
size/1,
handle_message/2
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
%% State
-record(state, {
local_node_id :: binary(),
routing_table :: macula_routing_table:routing_table(),
storage :: #{binary() => term()}, % Local key-value storage
config :: #{
k => pos_integer(),
alpha => pos_integer()
}
}).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start routing server with registered name macula_routing_server.
-spec start_link(binary(), map()) -> {ok, pid()} | {error, term()}.
start_link(LocalNodeId, Config) ->
gen_server:start_link({local, macula_routing_server}, ?MODULE, {LocalNodeId, Config}, []).
%% @doc Add node to routing table (async - does not block caller).
-spec add_node(pid(), macula_routing_bucket:node_info()) -> ok.
add_node(Pid, NodeInfo) ->
gen_server:cast(Pid, {add_node, NodeInfo}).
%% @doc Find k closest nodes to target.
-spec find_closest(pid(), binary(), pos_integer()) -> [macula_routing_bucket:node_info()].
find_closest(Pid, Target, K) ->
gen_server:call(Pid, {find_closest, Target, K}).
%% @doc Store value locally.
-spec store_local(pid(), binary(), term()) -> ok.
store_local(Pid, Key, Value) ->
gen_server:call(Pid, {store_local, Key, Value}).
%% @doc Store value in DHT by propagating to k closest nodes.
%% Stores locally first, then sends STORE messages to k closest peers.
-spec store(pid(), binary(), term()) -> ok.
store(Pid, Key, Value) ->
gen_server:call(Pid, {store, Key, Value}, 10000).
%% @doc Get value from local storage.
-spec get_local(pid(), binary()) -> {ok, term()} | not_found.
get_local(Pid, Key) ->
gen_server:call(Pid, {get_local, Key}).
%% @doc Get all keys from local storage.
-spec get_all_keys(pid()) -> {ok, [binary()]} | {error, term()}.
get_all_keys(Pid) ->
gen_server:call(Pid, get_all_keys).
%% @doc Find value in DHT using iterative lookup.
%% Returns {ok, Value} if found, {nodes, Nodes} if not found.
-spec find_value(pid(), binary(), pos_integer()) ->
{ok, term()} | {nodes, [macula_routing_bucket:node_info()]} | {error, term()}.
find_value(Pid, Key, K) ->
gen_server:call(Pid, {find_value, Key, K}, 10000).
%% @doc Get routing table snapshot.
-spec get_routing_table(pid()) -> macula_routing_table:routing_table().
get_routing_table(Pid) ->
gen_server:call(Pid, get_routing_table).
%% @doc Get number of nodes in routing table.
-spec size(pid()) -> non_neg_integer().
size(Pid) ->
gen_server:call(Pid, size).
%% @doc Handle incoming DHT message and return reply.
-spec handle_message(pid(), map()) -> map().
handle_message(Pid, Message) ->
gen_server:call(Pid, {handle_message, Message}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({LocalNodeId, Config}) ->
K = maps:get(k, Config, 20),
Alpha = maps:get(alpha, Config, 3),
State = #state{
local_node_id = LocalNodeId,
routing_table = macula_routing_table:new(LocalNodeId, K),
storage = #{},
config = #{
k => K,
alpha => Alpha
}
},
{ok, State}.
%% @private
handle_call({find_closest, Target, K}, _From, #state{routing_table = Table} = State) ->
Closest = macula_routing_table:find_closest(Table, Target, K),
{reply, Closest, State};
handle_call({store_local, Key, Value}, _From, #state{storage = Storage} = State) ->
ExistingProviders = maps:get(Key, Storage, []),
ProviderList = ensure_provider_list(ExistingProviders),
NodeId = maps:get(node_id, Value, undefined),
UpdatedProviders = upsert_provider(NodeId, Value, ProviderList),
NewStorage = Storage#{Key => UpdatedProviders},
{reply, ok, State#state{storage = NewStorage}};
handle_call({store, Key, Value}, _From, #state{routing_table = Table, config = Config} = State) ->
%% 1. Store locally first
NewState = case handle_call({store_local, Key, Value}, _From, State) of
{reply, ok, S} -> S;
_ -> State %% Shouldn't happen but be safe
end,
%% 2. Find k closest nodes to Key
K = maps:get(k, Config, 20),
ClosestNodes = macula_routing_table:find_closest(Table, Key, K),
%% 3. Send STORE message to each node (fire-and-forget)
StoreMsg = macula_routing_protocol:encode_store(Key, Value),
case ClosestNodes of
[] ->
%% No nodes in routing table - this is an embedded gateway connected to bootstrap
%% Forward the STORE to bootstrap gateway via RPC instead
forward_store_to_bootstrap(Key, Value);
_ ->
lists:foreach(fun(NodeInfo) ->
%% Best effort - don't fail if send fails
_ = macula_gateway_dht:send_to_peer(NodeInfo, store, StoreMsg)
end, ClosestNodes)
end,
{reply, ok, NewState};
handle_call({get_local, Key}, _From, #state{storage = Storage} = State) ->
Reply = format_storage_value(maps:get(Key, Storage, undefined)),
{reply, Reply, State};
handle_call(get_all_keys, _From, #state{storage = Storage} = State) ->
Keys = maps:keys(Storage),
{reply, {ok, Keys}, State};
handle_call(get_routing_table, _From, #state{routing_table = Table} = State) ->
{reply, Table, State};
handle_call(size, _From, #state{routing_table = Table} = State) ->
Size = macula_routing_table:size(Table),
{reply, Size, State};
handle_call({delete_local, Key, NodeId}, _From, #state{storage = Storage} = State) ->
NewStorage = delete_provider_from_storage(Key, NodeId, Storage),
{reply, ok, State#state{storage = NewStorage}};
handle_call({find_value, Key, K}, _From, #state{routing_table = Table, storage = Storage} = State) ->
?LOG_DEBUG("find_value: key=~p, storage_size=~p", [Key, maps:size(Storage)]),
?LOG_DEBUG("find_value: storage keys=~p", [maps:keys(Storage)]),
Reply = find_value_in_storage(Key, K, Storage, Table),
{reply, Reply, State};
handle_call({handle_message, Message}, _From, State) ->
{Reply, NewState} = process_dht_message(Message, State),
{reply, Reply, NewState};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
%% Add node to routing table (async operation)
handle_cast({add_node, NodeInfo}, #state{routing_table = Table} = State) ->
NewTable = macula_routing_table:add_node(Table, NodeInfo),
{noreply, State#state{routing_table = NewTable}};
%% @private
%% Handle async store (from _dht.store RPC) - stores locally only without propagation
%% This is used by bootstrap gateway when receiving store requests from peers
handle_cast({store, Key, Value}, #state{storage = Storage} = State) when is_binary(Key) ->
?LOG_DEBUG("Async STORE received - key hash prefix: ~p, value: ~p",
[binary:part(Key, 0, min(8, byte_size(Key))), Value]),
ExistingProviders = maps:get(Key, Storage, []),
ProviderList = ensure_provider_list(ExistingProviders),
NodeId = get_node_id_from_value(Value),
UpdatedProviders = upsert_provider(NodeId, Value, ProviderList),
NewStorage = Storage#{Key => UpdatedProviders},
?LOG_DEBUG("Stored value for key, now have ~p provider(s)", [length(UpdatedProviders)]),
{noreply, State#state{storage = NewStorage}};
handle_cast({store, Key, _Value}, State) ->
?LOG_ERROR("Async STORE received with non-binary key: ~p", [Key]),
{noreply, State};
handle_cast(_Request, State) ->
{noreply, State}.
%% @private
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Forward DHT STORE to bootstrap gateway via RPC when routing table is empty.
%% This allows embedded gateways to propagate subscriptions to the central bootstrap DHT.
-spec forward_store_to_bootstrap(binary(), term()) -> ok.
forward_store_to_bootstrap(Key, Value) ->
%% Call the _dht.store RPC procedure on the bootstrap gateway
%% This is a fire-and-forget operation - we don't wait for response
case whereis(macula_rpc_handler) of
undefined ->
ok;
_RpcHandler ->
%% Call _dht.store(Key, Value) on the bootstrap gateway
%% Using cast (fire-and-forget) since we don't need the response
try
Procedure = <<"_dht.store">>,
Args = #{
<<"key">> => Key,
<<"value">> => Value
},
%% Use macula module's call function which handles bootstrap routing
spawn(fun() ->
case whereis(macula_local_client) of
undefined ->
ok;
LocalClient ->
_ = macula:call(LocalClient, Procedure, Args, #{timeout => 5000}),
ok
end
end),
ok
catch
_:_Error ->
ok
end
end.
%% @doc Extract node_id from value map, handling both atom and binary keys.
-spec get_node_id_from_value(map()) -> binary() | undefined.
get_node_id_from_value(#{node_id := NodeId}) -> NodeId;
get_node_id_from_value(#{<<"node_id">> := NodeId}) -> NodeId;
get_node_id_from_value(_Value) -> undefined.
%% @doc Find index of provider with matching node_id in provider list.
-spec find_provider_index(binary(), [map()]) -> pos_integer() | not_found.
find_provider_index(NodeId, ProviderList) ->
find_provider_index(NodeId, ProviderList, 1).
find_provider_index(_NodeId, [], _Index) ->
not_found;
find_provider_index(NodeId, [Provider | Rest], Index) ->
ProviderNodeId = get_node_id_from_value(Provider),
check_provider_match(NodeId, ProviderNodeId, Rest, Index).
check_provider_match(NodeId, NodeId, _Rest, Index) -> Index;
check_provider_match(NodeId, _Other, Rest, Index) ->
find_provider_index(NodeId, Rest, Index + 1).
%% @doc Ensure value is a list for multi-provider storage.
-spec ensure_provider_list(term()) -> [term()].
ensure_provider_list(List) when is_list(List) -> List;
ensure_provider_list(Value) -> [Value].
%% @doc Insert or update a provider in the list.
-spec upsert_provider(binary() | undefined, map(), [map()]) -> [map()].
upsert_provider(undefined, Value, ProviderList) ->
[Value | ProviderList];
upsert_provider(NodeId, Value, ProviderList) ->
upsert_by_index(find_provider_index(NodeId, ProviderList), Value, ProviderList).
upsert_by_index(not_found, Value, ProviderList) ->
[Value | ProviderList];
upsert_by_index(Index, Value, ProviderList) ->
lists:sublist(ProviderList, Index - 1) ++ [Value] ++ lists:nthtail(Index, ProviderList).
%% @doc Format storage value for get_local response.
-spec format_storage_value(undefined | term()) -> not_found | {ok, [term()]}.
format_storage_value(undefined) -> not_found;
format_storage_value(Value) when is_list(Value) -> {ok, Value};
format_storage_value(Value) -> {ok, [Value]}.
%% @doc Delete a specific provider from storage by node_id.
-spec delete_provider_from_storage(binary(), binary(), map()) -> map().
delete_provider_from_storage(Key, NodeId, Storage) ->
case maps:get(Key, Storage, undefined) of
undefined -> Storage;
Providers when is_list(Providers) ->
delete_provider_by_node_id(Key, NodeId, Providers, Storage);
_SingleValue -> maps:remove(Key, Storage)
end.
delete_provider_by_node_id(Key, NodeId, Providers, Storage) ->
FilterFn = fun(P) -> get_node_id_from_value(P) =/= NodeId end,
UpdatedProviders = lists:filter(FilterFn, Providers),
update_or_remove_key(Key, UpdatedProviders, Storage).
update_or_remove_key(Key, [], Storage) -> maps:remove(Key, Storage);
update_or_remove_key(Key, Providers, Storage) -> Storage#{Key => Providers}.
%% @doc Find value in local storage or DHT.
-spec find_value_in_storage(binary(), pos_integer(), map(), macula_routing_table:routing_table()) ->
{ok, term()} | {error, term()}.
find_value_in_storage(Key, K, Storage, Table) ->
case maps:get(Key, Storage, undefined) of
undefined -> find_value_via_dht(Key, K, Table);
Value when is_list(Value) -> {ok, Value};
Value -> {ok, [Value]}
end.
find_value_via_dht(Key, K, Table) ->
QueryFn = fun(_NodeInfo, _QueryKey) -> {nodes, [_NodeInfo]} end,
case macula_routing_dht:find_value(Table, Key, K, QueryFn) of
{ok, Value} -> {ok, Value};
{nodes, _Nodes} -> {ok, []};
{error, Reason} -> {error, Reason}
end.
%% @doc Process incoming DHT message and generate reply.
-spec process_dht_message(map(), #state{}) -> {map(), #state{}}.
process_dht_message(Message, State) ->
MessageType = classify_dht_message(Message),
dispatch_dht_message(MessageType, Message, State).
classify_dht_message(Message) ->
case macula_routing_protocol:is_find_node(Message) of
true -> find_node;
false -> classify_non_find_node(Message)
end.
classify_non_find_node(Message) ->
case macula_routing_protocol:is_store(Message) of
true -> store;
false -> classify_find_value(Message)
end.
classify_find_value(Message) ->
case macula_routing_protocol:is_find_value(Message) of
true -> find_value;
false -> unknown
end.
dispatch_dht_message(find_node, Message, State) ->
handle_find_node(Message, State);
dispatch_dht_message(store, Message, State) ->
handle_store(Message, State);
dispatch_dht_message(find_value, Message, State) ->
handle_find_value(Message, State);
dispatch_dht_message(unknown, _Message, State) ->
{#{type => error, reason => unknown_message}, State}.
%% @doc Handle FIND_NODE request.
-spec handle_find_node(map(), #state{}) -> {map(), #state{}}.
handle_find_node(Message, #state{routing_table = Table, config = Config} = State) ->
{ok, Target} = macula_routing_protocol:decode_find_node(Message),
K = maps:get(k, Config, 20),
%% Find k closest nodes
Closest = macula_routing_table:find_closest(Table, Target, K),
%% Encode reply
Reply = macula_routing_protocol:encode_find_node_reply(Closest),
{Reply, State}.
%% @doc Handle STORE request.
-spec handle_store(map(), #state{}) -> {map(), #state{}}.
handle_store(Message, #state{storage = Storage} = State) ->
{ok, Key, Value} = macula_routing_protocol:decode_store(Message),
?LOG_DEBUG("STORE: key=~p, value=~p", [Key, Value]),
ExistingProviders = maps:get(Key, Storage, []),
ProviderList = ensure_provider_list(ExistingProviders),
NodeId = get_node_id_from_value(Value),
UpdatedProviders = upsert_provider(NodeId, Value, ProviderList),
NewStorage = Storage#{Key => UpdatedProviders},
?LOG_DEBUG("STORE complete: new_storage_size=~p, key_prefix=~p",
[maps:size(NewStorage), binary:part(Key, 0, min(8, byte_size(Key)))]),
Reply = #{type => store_reply, result => ok},
{Reply, State#state{storage = NewStorage}}.
%% @doc Handle FIND_VALUE request.
-spec handle_find_value(map(), #state{}) -> {map(), #state{}}.
handle_find_value(Message, #state{storage = Storage, routing_table = Table, config = Config} = State) ->
{ok, Key} = macula_routing_protocol:decode_find_value(Message),
StorageValue = maps:get(Key, Storage, undefined),
K = maps:get(k, Config, 20),
Reply = encode_find_value_result(StorageValue, Key, K, Table),
{Reply, State}.
encode_find_value_result(undefined, Key, K, Table) ->
Closest = macula_routing_table:find_closest(Table, Key, K),
macula_routing_protocol:encode_find_value_reply({nodes, Closest});
encode_find_value_result(Value, _Key, _K, _Table) when is_list(Value) ->
macula_routing_protocol:encode_find_value_reply({value, Value});
encode_find_value_result(Value, _Key, _K, _Table) ->
macula_routing_protocol:encode_find_value_reply({value, Value}).