Packages
macula
0.45.2
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_network.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Network transport and escalation operations for DHT routing.
%%% Extracted from macula_routing_server to reduce module size.
%%% Handles peer-to-peer store propagation, network queries,
%%% bridge escalation, and value replication.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_routing_network).
-include_lib("kernel/include/logger.hrl").
%% Store propagation
-export([
propagate_store_to_peers/4,
send_store_to_peer/2,
replicate_to_new_peer/4
]).
%% Network queries
-export([
network_query_find_value/2,
find_value_via_dht/3,
find_value_with_escalation/5
]).
%% Bridge/escalation
-export([
check_bridge_cache/1,
maybe_escalate_query/2,
escalate_to_bridge/2
]).
%%%===================================================================
%%% Store Propagation
%%%===================================================================
%% @doc Propagate STORE message to peers (called in spawned process).
-spec propagate_store_to_peers(list(), map(), binary(), term()) -> ok.
propagate_store_to_peers([], _StoreMsg, Key, Value) ->
?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.
%% @doc Send store message to a peer (best effort).
-spec send_store_to_peer(map(), map()) -> ok.
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.
%% @doc Replicate stored values to a newly joined peer.
-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 = macula_routing_storage: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, macula_routing_storage:format_node_id(NewNodeId)]);
false ->
ok
end,
ok.
%%%===================================================================
%%% Network Queries
%%%===================================================================
%% @doc Network query function for FIND_VALUE.
-spec network_query_find_value(map(), binary()) -> {value, term()} | {nodes, list()}.
network_query_find_value(NodeInfo, Key) ->
Endpoint = macula_routing_storage:get_node_endpoint(NodeInfo),
do_network_query_find_value(Endpoint, NodeInfo, Key).
%% @doc Query the DHT network for a value.
-spec find_value_via_dht(binary(), pos_integer(), macula_routing_table:routing_table()) ->
{ok, term()} | {error, term()}.
find_value_via_dht(Key, K, Table) ->
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), macula_routing_storage: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.
%% @doc Find value with escalation support.
-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) ->
case check_bridge_cache(Key) of
{ok, CachedValue} ->
?LOG_DEBUG("find_value: cache hit for key ~p", [Key]),
{ok, CachedValue};
not_found ->
LocalValues = case maps:get(Key, Storage, undefined) of
undefined -> [];
Value when is_list(Value) -> Value;
Value -> [Value]
end,
find_value_with_local(LocalValues, Key, K, Table, Config)
end.
%%%===================================================================
%%% Bridge / Escalation
%%%===================================================================
%% @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) ->
case maps:get(escalation_enabled, Config, true) 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.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @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]).
%% @private Forward DHT STORE to bootstrap gateway via RPC.
forward_store_to_bootstrap(Key, Value) ->
RpcHandler = whereis(macula_rpc_handler),
do_forward_store(RpcHandler, Key, Value).
do_forward_store(undefined, _Key, _Value) ->
ok;
do_forward_store(_RpcHandler, Key, Value) ->
Procedure = <<"_dht.store">>,
Args = #{<<"key">> => Key, <<"value">> => Value},
spawn(fun() -> forward_store_via_local_client(Procedure, Args) end),
ok.
forward_store_via_local_client(Procedure, Args) ->
LocalClient = whereis(macula_local_client),
do_forward_via_client(LocalClient, Procedure, Args).
do_forward_via_client(undefined, _Procedure, _Args) ->
ok;
do_forward_via_client(LocalClient, Procedure, Args) ->
handle_forward_result(catch macula:call(LocalClient, Procedure, Args, #{timeout => 5000})).
handle_forward_result({'EXIT', _}) -> ok;
handle_forward_result(_Result) -> ok.
%% @private
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) ->
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),
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.
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 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) ->
?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
[] -> maybe_escalate_query(Key, Config);
_ -> {ok, NetworkValues}
end.