Packages
macula
0.38.8
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_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,
remove_node/2,
find_closest/3,
store_local/3,
store/3,
get_local/2,
get_all_keys/1,
delete_local/3,
find_value/3,
find_value_local/2,
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
]).
%% Stale peer eviction interval (60 seconds)
-define(EVICT_INTERVAL_MS, 60000).
%% Peers not seen in 5 minutes are considered stale
-define(STALE_THRESHOLD_MS, 300000).
%% DHT value expiry interval (30 seconds)
-define(VALUE_EXPIRY_INTERVAL_MS, 30000).
%% DHT values with TTL older than this are expired (default 5 minutes + 60s grace)
-define(VALUE_MAX_AGE_MS, 360000).
%% 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(),
escalation_enabled => boolean(),
escalation_timeout => 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 Remove node from routing table (async - does not block caller).
-spec remove_node(pid(), binary()) -> ok.
remove_node(Pid, NodeId) ->
gen_server:cast(Pid, {remove_node, NodeId}).
%% @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 — goes through gen_server for full iterative lookup.
%% For local-only reads (answering peer queries), use find_value_local/2.
-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 Fast local-only lookup via ETS — for answering incoming FIND_VALUE from peers.
%% Does NOT query the network. Returns only locally stored values.
-spec find_value_local(binary(), pos_integer()) ->
{ok, term()} | {error, not_found}.
find_value_local(Key, _K) ->
case ets:info(macula_dht_storage) of
undefined -> {error, not_found};
_ ->
case ets:lookup(macula_dht_storage, Key) of
[{_, Values}] when is_list(Values), Values =/= [] -> {ok, Values};
_ -> {error, not_found}
end
end.
%% @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),
EscalationEnabled = maps:get(escalation_enabled, Config, true),
EscalationTimeout = maps:get(escalation_timeout, Config, 5000),
%% Named ETS table for concurrent DHT storage reads (bypasses gen_server mailbox)
StorageEts = ets:new(macula_dht_storage, [named_table, set, public, {read_concurrency, true}]),
State = #state{
local_node_id = LocalNodeId,
routing_table = macula_routing_table:new(LocalNodeId, K),
storage = #{},
config = #{
k => K,
alpha => Alpha,
escalation_enabled => EscalationEnabled,
escalation_timeout => EscalationTimeout,
storage_ets => StorageEts
}
},
%% Start periodic stale peer eviction
erlang:send_after(?EVICT_INTERVAL_MS, self(), evict_stale_peers),
%% Start periodic DHT value expiry (removes zombie subscriptions)
erlang:send_after(?VALUE_EXPIRY_INTERVAL_MS, self(), expire_stale_values),
{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, config = Config} = 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},
%% Mirror to ETS for concurrent reads
case maps:get(storage_ets, Config, undefined) of
undefined -> ok;
Ets -> ets:insert(Ets, {Key, UpdatedProviders})
end,
{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, config = Config} = State) ->
NewStorage = delete_provider_from_storage(Key, NodeId, Storage),
%% Mirror to ETS
case maps:get(storage_ets, Config, undefined) of
undefined -> ok;
Ets ->
case maps:get(Key, NewStorage, []) of
[] -> ets:delete(Ets, Key);
V -> ets:insert(Ets, {Key, V})
end
end,
{reply, ok, State#state{storage = NewStorage}};
handle_call({find_value, Key, K}, From, #state{routing_table = Table, storage = Storage,
config = Config} = State) ->
?LOG_DEBUG("find_value: key=~p, storage_size=~p", [Key, maps:size(Storage)]),
LocalValues = case maps:get(Key, Storage, undefined) of
undefined -> [];
V when is_list(V) -> V;
V -> [V]
end,
%% Always query the network — local values may be incomplete (e.g. pubsub subscribers).
%% Spawn to avoid blocking the routing server.
spawn(fun() ->
RemoteResult = find_value_with_escalation(Key, K, Storage, Table, Config),
RemoteValues = case RemoteResult of
{ok, R} when is_list(R) -> R;
_ -> []
end,
%% Merge local + remote, deduplicate by node_id
AllValues = deduplicate_providers(LocalValues ++ RemoteValues),
Reply = case AllValues of
[] -> {error, not_found};
_ -> {ok, AllValues}
end,
gen_server:reply(From, Reply)
end),
{noreply, 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)
%% When a new node actually joins, replicate relevant STORE values to it.
%% This is standard Kademlia: new neighbors receive stored values they are
%% now among the k-closest nodes for.
handle_cast({add_node, NodeInfo}, #state{routing_table = Table, storage = Storage,
config = Config} = State) ->
OldSize = macula_routing_table:size(Table),
NewTable = macula_routing_table:add_node(Table, NodeInfo),
NewSize = macula_routing_table:size(NewTable),
NodeId = maps:get(node_id, NodeInfo, undefined),
?LOG_DEBUG("[RoutingServer] add_node: node_id=~s, table_size: ~p -> ~p",
[case NodeId of B when is_binary(B), byte_size(B) =:= 32 -> binary:encode_hex(B);
B when is_binary(B) -> B; _ -> <<"?">> end,
OldSize, NewSize]),
%% Replicate stored values to the new peer if it actually joined
case NewSize > OldSize of
true ->
spawn(fun() -> replicate_to_new_peer(NodeInfo, Storage, NewTable, Config) end);
false ->
ok
end,
{noreply, State#state{routing_table = NewTable}};
handle_cast({remove_node, NodeId}, #state{routing_table = Table} = State) ->
OldSize = macula_routing_table:size(Table),
NewTable = macula_routing_table:remove_node(Table, NodeId),
NewSize = macula_routing_table:size(NewTable),
?LOG_DEBUG("[RoutingServer] remove_node: node_id=~s, table_size: ~p -> ~p",
[binary:encode_hex(NodeId), OldSize, NewSize]),
{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, config = Config} = 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},
%% Mirror to ETS for concurrent reads
case maps:get(storage_ets, Config, undefined) of
undefined -> ok;
Ets -> ets:insert(Ets, {Key, UpdatedProviders})
end,
?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(evict_stale_peers, #state{routing_table = Table} = State) ->
StaleThreshold = erlang:system_time(millisecond) - ?STALE_THRESHOLD_MS,
OldSize = macula_routing_table:size(Table),
NewTable = macula_routing_table:evict_stale(Table, StaleThreshold),
NewSize = macula_routing_table:size(NewTable),
case OldSize - NewSize of
0 -> ok;
Evicted ->
?LOG_INFO("[RoutingServer] Evicted ~p stale peer(s) (table: ~p -> ~p)",
[Evicted, OldSize, NewSize])
end,
erlang:send_after(?EVICT_INTERVAL_MS, self(), evict_stale_peers),
{noreply, State#state{routing_table = NewTable}};
handle_info(expire_stale_values, #state{storage = Storage, config = Config} = State) ->
Now = erlang:system_time(millisecond),
{NewStorage, Expired} = maps:fold(fun(Key, Providers, {StorageAcc, ExpiredCount}) ->
Fresh = [P || P <- Providers, not is_provider_expired(P, Now)],
case {Fresh, length(Providers) - length(Fresh)} of
{[], Removed} ->
{maps:remove(Key, StorageAcc), ExpiredCount + Removed};
{Kept, 0} ->
{StorageAcc#{Key => Kept}, ExpiredCount};
{Kept, Removed} ->
{StorageAcc#{Key => Kept}, ExpiredCount + Removed}
end
end, {Storage, 0}, Storage),
case Expired of
0 -> ok;
N -> ?LOG_INFO("[RoutingServer] Expired ~p stale DHT value(s)", [N])
end,
%% Mirror to ETS
case Expired > 0 of
true ->
case maps:get(storage_ets, Config, undefined) of
undefined -> ok;
Ets ->
ets:delete_all_objects(Ets),
maps:foreach(fun(K, V) -> ets:insert(Ets, {K, V}) end, NewStorage)
end;
false -> ok
end,
erlang:send_after(?VALUE_EXPIRY_INTERVAL_MS, self(), expire_stale_values),
{noreply, State#state{storage = NewStorage}};
handle_info(_Info, State) ->
{noreply, State}.
%% @private Check if a provider entry is expired based on its stored_at timestamp.
%% Providers without stored_at are considered expired if they have a ttl field
%% (indicating they're subscription entries that should have been timestamped).
-spec is_provider_expired(map(), integer()) -> boolean().
is_provider_expired(#{stored_at := StoredAt}, Now) ->
(Now - StoredAt) > ?VALUE_MAX_AGE_MS;
is_provider_expired(#{<<"stored_at">> := StoredAt}, Now) ->
%% Binary key variant (from msgpack deserialization over network)
(Now - StoredAt) > ?VALUE_MAX_AGE_MS;
is_provider_expired(#{ttl := _TTL}, _Now) ->
%% Has TTL but no stored_at — legacy entry, expire it
true;
is_provider_expired(#{<<"ttl">> := _TTL}, _Now) ->
%% Binary key variant — legacy entry, expire it
true;
is_provider_expired(_Provider, _Now) ->
%% No TTL, no stored_at — service entry, keep it
false.
%% @private
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Replicate stored values to a newly joined peer (called in spawned process).
%% Standard Kademlia behavior: when a new node joins the routing table,
%% send it any stored values for which it is now among the k-closest nodes.
-spec replicate_to_new_peer(map(), map(), macula_routing_table:routing_table(), map()) -> ok.
replicate_to_new_peer(NewNodeInfo, Storage, Table, Config) ->
K = maps:get(k, Config, 20),
NewNodeId = maps:get(node_id, NewNodeInfo, undefined),
Replicated = maps:fold(fun(Key, Providers, Acc) ->
ClosestNodes = macula_routing_table:find_closest(Table, Key, K),
ClosestIds = [maps:get(node_id, N, undefined) || N <- ClosestNodes],
case lists:member(NewNodeId, ClosestIds) of
true ->
ProviderList = ensure_provider_list(Providers),
lists:foreach(fun(Provider) ->
StoreMsg = macula_routing_protocol:encode_store(Key, Provider),
send_store_to_peer(NewNodeInfo, StoreMsg)
end, ProviderList),
Acc + length(ProviderList);
false ->
Acc
end
end, 0, Storage),
case Replicated > 0 of
true ->
?LOG_DEBUG("[RoutingServer] Replicated ~p stored value(s) to new peer ~s",
[Replicated, format_node_id(NewNodeId)]);
false ->
ok
end,
ok.
%% @doc Format node_id for logging.
-spec format_node_id(binary() | undefined) -> binary().
format_node_id(undefined) -> <<"?">>;
format_node_id(B) when is_binary(B), byte_size(B) =:= 32 -> binary:encode_hex(B);
format_node_id(B) when is_binary(B) -> B;
format_node_id(_) -> <<"?">>.
%% @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
?LOG_DEBUG("[DHT] No peers in routing table, forwarding store to bootstrap"),
forward_store_to_bootstrap(Key, Value);
propagate_store_to_peers(ClosestNodes, StoreMsg, _Key, _Value) ->
?LOG_DEBUG("[DHT] Propagating store to ~p peer(s)", [length(ClosestNodes)]),
lists:foreach(fun(NodeInfo) ->
send_store_to_peer(NodeInfo, StoreMsg)
end, ClosestNodes),
ok.
%% @private Send store message to a peer (best effort)
send_store_to_peer(NodeInfo, StoreMsg) ->
?LOG_DEBUG("[DHT] Sending store to peer ~p", [NodeInfo]),
Transport = persistent_term:get(macula_dht_transport, macula_gateway_dht),
try Transport:send_to_peer(NodeInfo, store, StoreMsg) of
Result -> handle_peer_send_result(Result, NodeInfo)
catch
error:function_clause:Stacktrace ->
?LOG_WARNING("[DHT] Store send function_clause to ~p:~n ~p", [NodeInfo, Stacktrace]);
Class:Error:Stacktrace ->
?LOG_WARNING("[DHT] Store send failed to ~p: ~p:~p~n ~p",
[NodeInfo, Class, Error, Stacktrace])
end.
%% @private Handle peer send result
handle_peer_send_result({error, Reason}, NodeInfo) ->
?LOG_DEBUG("[DHT] Store send error to ~p: ~p", [NodeInfo, Reason]);
handle_peer_send_result(Result, _NodeInfo) ->
?LOG_DEBUG("[DHT] Store send result: ~p", [Result]).
%% @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
RpcHandler = whereis(macula_rpc_handler),
do_forward_store(RpcHandler, Key, Value).
%% @private RPC handler not available
do_forward_store(undefined, _Key, _Value) ->
ok;
%% @private RPC handler available - spawn and forward
do_forward_store(_RpcHandler, Key, Value) ->
Procedure = <<"_dht.store">>,
Args = #{
<<"key">> => Key,
<<"value">> => Value
},
%% Spawn fire-and-forget call to bootstrap
spawn(fun() -> forward_store_via_local_client(Procedure, Args) end),
ok.
%% @private Forward store via local client (spawned process)
forward_store_via_local_client(Procedure, Args) ->
LocalClient = whereis(macula_local_client),
do_forward_via_client(LocalClient, Procedure, Args).
%% @private Local client not available
do_forward_via_client(undefined, _Procedure, _Args) ->
ok;
%% @private Local client available - make call
do_forward_via_client(LocalClient, Procedure, Args) ->
handle_forward_result(catch macula:call(LocalClient, Procedure, Args, #{timeout => 5000})).
%% @private Handle forward result (fire-and-forget, ignore errors)
handle_forward_result({'EXIT', _}) ->
ok;
handle_forward_result(_Result) ->
ok.
%% @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.
%% Stamps every entry with stored_at for expiry tracking.
-spec upsert_provider(binary() | undefined, map(), [map()]) -> [map()].
upsert_provider(undefined, Value, ProviderList) ->
[Value#{stored_at => erlang:system_time(millisecond)} | ProviderList];
upsert_provider(NodeId, Value, ProviderList) ->
Stamped = Value#{stored_at => erlang:system_time(millisecond)},
upsert_by_index(find_provider_index(NodeId, ProviderList), Stamped, 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 Deduplicate providers by node_id — keeps first occurrence.
deduplicate_providers(Providers) ->
lists:foldl(fun(P, Acc) ->
NId = get_node_id_from_value(P),
case lists:any(fun(A) -> get_node_id_from_value(A) =:= NId end, Acc) of
true -> Acc;
false -> [P | Acc]
end
end, [], Providers).
%% @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 with escalation support.
%% Local-first lookup: returns local values immediately if found.
%% Only queries DHT network if no local values (for RPC provider discovery).
%% For pub/sub, bootstrap stores all subscriptions locally - no network query needed.
-spec find_value_with_escalation(binary(), pos_integer(), map(),
macula_routing_table:routing_table(), map()) ->
{ok, term()} | {error, term()}.
find_value_with_escalation(Key, K, Storage, Table, Config) ->
%% 1. Check bridge cache first (fastest)
case check_bridge_cache(Key) of
{ok, CachedValue} ->
?LOG_DEBUG("find_value: cache hit for key ~p", [Key]),
{ok, CachedValue};
not_found ->
%% 2. Get local values (if any)
LocalValues = case maps:get(Key, Storage, undefined) of
undefined -> [];
Value when is_list(Value) -> Value;
Value -> [Value]
end,
%% 3. If local values found, return immediately (local-first)
%% For pub/sub with bootstrap pattern, all subscriptions are local
find_value_with_local(LocalValues, Key, K, Table, Config)
end.
%% @private Return local values immediately if found, else query network
find_value_with_local(LocalValues, _Key, _K, _Table, _Config) when LocalValues =/= [] ->
?LOG_DEBUG("[DHT] find_value: returning ~p local value(s)", [length(LocalValues)]),
{ok, LocalValues};
find_value_with_local([], Key, K, Table, Config) ->
%% No local values - query network for remote providers (RPC discovery)
?LOG_DEBUG("[DHT] find_value: no local values, querying network"),
NetworkValues = case find_value_via_dht(Key, K, Table) of
{ok, RemoteList} when is_list(RemoteList) -> RemoteList;
{ok, RemoteSingle} -> [RemoteSingle];
_ -> []
end,
?LOG_DEBUG("[DHT] find_value: network returned ~p value(s)", [length(NetworkValues)]),
case NetworkValues of
[] ->
%% Try escalation if nothing found
maybe_escalate_query(Key, Config);
_ ->
{ok, NetworkValues}
end.
find_value_via_dht(Key, K, Table) ->
%% Log routing table state for debugging
InitialClosest = macula_routing_table:find_closest(Table, Key, K),
?LOG_DEBUG("[DHT] find_value_via_dht: routing_table has ~p nodes for key lookup", [length(InitialClosest)]),
lists:foreach(fun(NodeInfo) ->
?LOG_DEBUG("[DHT] find_value_via_dht: will query node ~p at ~p",
[maps:get(node_id, NodeInfo, unknown), get_node_endpoint(NodeInfo)])
end, InitialClosest),
QueryFn = fun(NodeInfo, QueryKey) -> network_query_find_value(NodeInfo, QueryKey) end,
case macula_routing_dht:find_value(Table, Key, K, QueryFn) of
{ok, Value} ->
?LOG_DEBUG("[DHT] find_value_via_dht: found ~p subscriber(s)", [length(Value)]),
{ok, Value};
{nodes, Nodes} ->
?LOG_DEBUG("[DHT] find_value_via_dht: not found, got ~p nodes", [length(Nodes)]),
{ok, []};
{error, Reason} ->
?LOG_WARNING("[DHT] find_value_via_dht: error ~p", [Reason]),
{error, Reason}
end.
%% @private Network query function for FIND_VALUE.
%% Sends FIND_VALUE request to a peer and waits for response.
%% Returns {value, Value} if peer has the value, {nodes, []} otherwise.
-spec network_query_find_value(map(), binary()) -> {value, term()} | {nodes, list()}.
network_query_find_value(NodeInfo, Key) ->
Endpoint = get_node_endpoint(NodeInfo),
do_network_query_find_value(Endpoint, NodeInfo, Key).
do_network_query_find_value(undefined, NodeInfo, _Key) ->
?LOG_WARNING("[DHT] network_query_find_value: no endpoint for node ~p", [NodeInfo]),
{nodes, []};
do_network_query_find_value(Endpoint, NodeInfo, Key) ->
%% Build FIND_VALUE message (the protocol expects raw Key, not encoded)
FindValueMsg = #{<<"key">> => Key},
?LOG_DEBUG("[DHT] network_query_find_value: querying ~s for key", [Endpoint]),
Transport = persistent_term:get(macula_dht_transport, macula_gateway_dht),
%% Send request and wait for response (synchronous with 5s timeout)
case Transport:send_and_wait(NodeInfo, find_value, FindValueMsg, 5000) of
{ok, {find_value_reply, Response}} ->
Result = decode_find_value_response(Response),
?LOG_DEBUG("[DHT] network_query_find_value: got reply from ~s: ~p", [Endpoint, Result]),
Result;
{ok, {OtherType, _Response}} ->
?LOG_WARNING("[DHT] network_query_find_value: unexpected reply type ~p from ~s", [OtherType, Endpoint]),
{nodes, []};
{error, timeout} ->
?LOG_WARNING("[DHT] network_query_find_value: timeout querying ~s", [Endpoint]),
{nodes, []};
{error, Reason} ->
?LOG_WARNING("[DHT] network_query_find_value: error ~p querying ~s", [Reason, Endpoint]),
{nodes, []}
end.
%% @private Decode FIND_VALUE reply content.
decode_find_value_response(Response) ->
case macula_routing_protocol:decode_find_value_reply(Response) of
{ok, {value, Value}} ->
{value, Value};
{ok, {nodes, Nodes}} ->
{nodes, Nodes};
{error, _Reason} ->
{nodes, []}
end.
%% @private Extract endpoint from node info.
get_node_endpoint(#{endpoint := Endpoint}) when is_binary(Endpoint) ->
Endpoint;
get_node_endpoint(#{address := {Host, Port}}) when is_integer(Port) ->
iolist_to_binary([format_host(Host), <<":">>, integer_to_binary(Port)]);
get_node_endpoint(#{<<"endpoint">> := Endpoint}) when is_binary(Endpoint) ->
Endpoint;
get_node_endpoint(#{address := Address}) when is_binary(Address) ->
%% Address is already a full URL (e.g., <<"https://host:port">>)
Address;
get_node_endpoint(#{<<"address">> := Address}) when is_binary(Address) ->
%% Address is already a full URL (binary key version)
Address;
get_node_endpoint(_) ->
undefined.
%% @private Format host for URL.
format_host({A, B, C, D}) when is_integer(A) ->
io_lib:format("~B.~B.~B.~B", [A, B, C, D]);
format_host(Host) when is_list(Host) ->
Host;
format_host(Host) when is_binary(Host) ->
binary_to_list(Host).
%% @doc Check bridge cache for value.
-spec check_bridge_cache(binary()) -> {ok, term()} | not_found.
check_bridge_cache(Key) ->
case whereis(macula_bridge_cache) of
undefined -> not_found;
CachePid ->
case macula_bridge_cache:get(CachePid, Key) of
{ok, Value} -> {ok, Value};
_ -> not_found
end
end.
%% @doc Escalate query to parent bridge if enabled.
-spec maybe_escalate_query(binary(), map()) -> {ok, term()} | {error, term()}.
maybe_escalate_query(Key, Config) ->
EscalationEnabled = maps:get(escalation_enabled, Config, true),
case EscalationEnabled of
false ->
{ok, []};
true ->
escalate_to_bridge(Key, Config)
end.
%% @doc Escalate query to parent bridge.
-spec escalate_to_bridge(binary(), map()) -> {ok, term()} | {error, term()}.
escalate_to_bridge(Key, Config) ->
case whereis(macula_bridge_node) of
undefined ->
?LOG_DEBUG("find_value: no bridge node available for escalation"),
{ok, []};
BridgePid ->
Timeout = maps:get(escalation_timeout, Config, 5000),
Query = #{type => find_value, key => Key},
?LOG_DEBUG("find_value: escalating query for key ~p to bridge", [Key]),
case macula_bridge_node:escalate_query(BridgePid, Query, Timeout) of
{ok, Value} ->
?LOG_DEBUG("find_value: escalation successful, got value"),
{ok, Value};
{error, not_connected} ->
?LOG_DEBUG("find_value: bridge not connected to parent"),
{ok, []};
{error, Reason} ->
?LOG_WARNING("find_value: escalation failed: ~p", [Reason]),
{ok, []}
end
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
TableSize = macula_routing_table:size(Table),
Closest = macula_routing_table:find_closest(Table, Target, K),
?LOG_DEBUG("[RoutingServer] FIND_NODE: target=~s, table_size=~p, found=~p",
[binary:encode_hex(Target), TableSize, length(Closest)]),
%% 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("[RoutingServer] STORE received: key_prefix=~p, value_node_id=~p",
[binary:part(Key, 0, min(8, byte_size(Key))), maps:get(node_id, Value, maps:get(<<"node_id">>, Value, unknown))]),
?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),
?LOG_DEBUG("[RoutingServer] FIND_VALUE: key_prefix=~p, storage_size=~p, found=~p",
[binary:part(Key, 0, min(8, byte_size(Key))), maps:size(Storage), StorageValue =/= undefined]),
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}).