Packages

macula

0.12.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
macula src macula_routing_system macula_routing_server.erl
Raw

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,
delete_local/3,
find_value/3,
get_routing_table/1,
size/1,
handle_message/2,
handle_message_async/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 Delete value from local storage.
-spec delete_local(pid(), binary(), binary()) -> ok.
delete_local(Pid, Key, NodeId) ->
gen_server:call(Pid, {delete_local, Key, NodeId}).
%% @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}).
%% @doc Handle incoming DHT message asynchronously (fire-and-forget).
%% Use for STORE messages where no reply is needed.
-spec handle_message_async(pid(), map()) -> ok.
handle_message_async(Pid, Message) ->
gen_server:cast(Pid, {handle_message_async, 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, non-blocking)
%% Spawn the sends to avoid blocking the routing server on network I/O
StoreMsg = macula_routing_protocol:encode_store(Key, Value),
spawn(fun() -> propagate_store_to_peers(ClosestNodes, StoreMsg, Key, Value) 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};
%% @private
%% Handle async DHT message processing (fire-and-forget)
%% Used by gateway to avoid blocking on STORE operations
handle_cast({handle_message_async, Message}, State) ->
{_Reply, NewState} = process_dht_message(Message, State),
{noreply, NewState};
handle_cast(_Request, State) ->
{noreply, State}.
%% @private
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Propagate STORE message to peers (called in spawned process).
%% Non-blocking: network I/O happens in separate process, won't block routing server.
-spec propagate_store_to_peers(list(), map(), binary(), term()) -> ok.
propagate_store_to_peers([], _StoreMsg, Key, Value) ->
%% No nodes in routing table - forward to bootstrap gateway via RPC
forward_store_to_bootstrap(Key, Value);
propagate_store_to_peers(ClosestNodes, StoreMsg, _Key, _Value) ->
lists:foreach(fun(NodeInfo) ->
%% Best effort - don't fail if send fails
try
_ = macula_gateway_dht:send_to_peer(NodeInfo, store, StoreMsg)
catch
_:_Error -> ok
end
end, ClosestNodes),
ok.
%% @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}).