Packages

macula

0.45.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
macula src macula_gateway_system macula_gateway_dht.erl
Raw

src/macula_gateway_system/macula_gateway_dht.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% DHT Query Handler Module - handles DHT message forwarding to routing server.
%%%
%%% Responsibilities:
%%% - Forward DHT STORE messages to routing server
%%% - Forward DHT FIND_VALUE messages to routing server, send encoded replies
%%% - Forward DHT FIND_NODE messages to routing server, send encoded replies
%%% - Handle DHT queries from process messages
%%% - Encode replies using protocol encoder
%%% - Handle errors gracefully
%%%
%%% Pattern: Stateless delegation module
%%% - No GenServer (no state to manage)
%%% - Pure functions forwarding to routing server
%%% - Consistent error handling ({ok, Result} | {error, Reason})
%%%
%%% Extracted from macula_gateway.erl (Phase 10)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_dht).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
handle_store/2,
handle_find_value/2,
handle_find_node/2,
handle_query/3,
lookup_value/1,
do_handle_find_value/1,
forward_publish_to_bootstrap/1,
send_to_peer/3,
send_and_wait/4,
query_peer/3
]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Handle DHT STORE message.
%% Forwards to routing server asynchronously (fire-and-forget, no reply needed).
%% Uses async handler to prevent blocking the gateway on DHT operations.
-spec handle_store(reference(), map()) -> ok.
handle_store(_Stream, StoreMsg) ->
%% Forward to routing server asynchronously (fire-and-forget)
%% STORE doesn't require a response, so we don't wait
?LOG_DEBUG("Processing STORE message, forwarding to routing_server (async)"),
macula_routing_server:handle_message_async(macula_routing_server, StoreMsg),
ok.
%% @doc Handle DHT FIND_VALUE message.
%% Extracts the key and performs local storage lookup, returning result over stream.
%% The message format from protocol decoder contains a binary key field.
-spec handle_find_value(reference(), map()) -> ok.
handle_find_value(Stream, FindValueMsg) ->
%% Extract key from message - supports both <<"key">> (protocol) and key (atom)
Key = maps:get(<<"key">>, FindValueMsg, maps:get(key, FindValueMsg, undefined)),
?LOG_INFO("[Gateway DHT] FIND_VALUE request for key_prefix=~p",
[case Key of
undefined -> undefined;
K when is_binary(K) -> binary:part(K, 0, min(8, byte_size(K)));
_ -> Key
end]),
Reply = do_handle_find_value(Key),
ReplyBinary = macula_protocol_encoder:encode(find_value_reply, Reply),
macula_quic:send(Stream, ReplyBinary),
ok.
%% @private Handle FIND_VALUE with undefined key
do_handle_find_value(undefined) ->
?LOG_WARNING("[Gateway DHT] FIND_VALUE with undefined key"),
#{type => error, reason => missing_key};
%% @private Handle FIND_VALUE with valid key.
%% Uses fast ETS local-only lookup (no gen_server, no network query).
do_handle_find_value(Key) when is_binary(Key) ->
case macula_routing_server:find_value_local(Key, 20) of
{error, not_found} -> macula_routing_protocol:encode_find_value_reply({nodes, []});
Result -> do_handle_find_value_result(Result)
end;
do_handle_find_value(Key) ->
?LOG_WARNING("[Gateway DHT] FIND_VALUE with non-binary key: ~p", [Key]),
#{type => error, reason => invalid_key}.
do_handle_find_value_result({ok, Values}) when is_list(Values), Values =/= [] ->
?LOG_INFO("[Gateway DHT] FIND_VALUE found ~p value(s)", [length(Values)]),
macula_routing_protocol:encode_find_value_reply({value, Values});
do_handle_find_value_result({ok, []}) ->
macula_routing_protocol:encode_find_value_reply({nodes, []});
do_handle_find_value_result({error, _}) ->
macula_routing_protocol:encode_find_value_reply({nodes, []}).
%% @doc Handle DHT FIND_NODE message.
%% The message is already decoded by the gateway — it's the payload without the type field.
%% We extract the target and query the routing table directly instead of going through
%% handle_message (which re-classifies and fails because the type field is stripped).
-spec handle_find_node(reference(), map()) -> ok.
handle_find_node(Stream, FindNodeMsg) ->
Target = maps:get(<<"target">>, FindNodeMsg, maps:get(target, FindNodeMsg, undefined)),
Closest = find_closest_nodes(Target),
Reply = macula_routing_protocol:encode_find_node_reply(Closest),
ReplyBinary = macula_protocol_encoder:encode(find_node_reply, Reply),
macula_quic:send(Stream, ReplyBinary),
ok.
find_closest_nodes(undefined) ->
[];
find_closest_nodes(Target) ->
macula_routing_server:find_closest(macula_routing_server, Target, 20).
%% @doc Handle DHT query from process message.
%% Decodes query, forwards to routing server, encodes reply, sends to requesting process.
%% Crashes on decode or routing failures - exposes protocol/DHT bugs.
-spec handle_query(pid(), atom(), binary()) -> ok.
handle_query(FromPid, _QueryType, QueryData) ->
%% Decode the query message (let it crash on decode errors)
{ok, {MessageType, Message}} = macula_protocol_decoder:decode(QueryData),
%% Forward to DHT routing server (let it crash on errors)
Reply = macula_routing_server:handle_message(macula_routing_server, Message),
%% Encode reply based on message type (let it crash on encoding errors)
ReplyData = encode_reply_by_type(MessageType, Reply),
%% Send reply back to requesting process
FromPid ! {dht_reply, ReplyData},
ok.
%% @doc Look up ALL providers for a key — merges local + network results.
%% Used by pubsub discovery where we need ALL subscribers, not just the first.
-spec lookup_value(binary()) -> {ok, list()} | {error, not_found}.
lookup_value(Key) ->
%% Get local results (fast, no gen_server)
LocalValues = case macula_routing_server:find_value_local(Key, 20) of
{ok, L} when is_list(L) -> L;
_ -> []
end,
%% Also query peers for their stored values
RemoteValues = case whereis(macula_routing_server) of
undefined -> [];
Pid ->
case safe_find_value(Pid, Key) of
{ok, R} when is_list(R) -> R;
_ -> []
end
end,
%% Merge and deduplicate by node_id
Merged = merge_providers(LocalValues ++ RemoteValues),
case Merged of
[] -> {error, not_found};
_ -> {ok, Merged}
end.
%% @private Deduplicate providers by node_id
merge_providers(Providers) ->
lists:foldl(fun(P, Acc) ->
NodeId = get_provider_node_id(P),
case lists:any(fun(A) -> get_provider_node_id(A) =:= NodeId end, Acc) of
true -> Acc;
false -> [P | Acc]
end
end, [], Providers).
get_provider_node_id(#{node_id := N}) -> N;
get_provider_node_id(#{<<"node_id">> := N}) -> N;
get_provider_node_id(_) -> make_ref().
%% @private Safe find_value via gen_server — catches timeouts
safe_find_value(Pid, Key) ->
try macula_routing_server:find_value(Pid, Key, 20)
catch
exit:{timeout, _} -> {error, timeout}
end.
%% @doc Forward a PUBLISH message to the bootstrap gateway for distribution.
%% @deprecated v0.14.0+ uses direct P2P routing via local DHT lookup.
%% Bootstrap is NOT a broker - use route_via_local_dht in pubsub_router instead.
%% This function remains for backwards compatibility but should not be used.
-spec forward_publish_to_bootstrap(map()) -> ok | {error, term()}.
forward_publish_to_bootstrap(PubMsg) ->
Realm = application:get_env(macula, realm, <<"default">>),
do_forward_to_bootstrap(gproc:lookup_local_name({connection, Realm}), PubMsg).
do_forward_to_bootstrap(undefined, _PubMsg) ->
{error, no_connection};
do_forward_to_bootstrap(ConnPid, PubMsg) ->
macula_connection:send_message(ConnPid, publish, PubMsg).
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Send a DHT message to a peer and wait for response.
%% Used by macula_routing_server for iterative FIND_VALUE lookups.
-spec send_and_wait(map(), atom(), map(), timeout()) -> {ok, term()} | {error, term()}.
send_and_wait(NodeInfo, MessageType, Message, Timeout) ->
Endpoint = extract_endpoint(NodeInfo),
case Endpoint of
undefined -> {error, no_endpoint};
_ -> macula_peer_connector:send_message_and_wait(Endpoint, MessageType, Message, Timeout)
end.
-spec send_to_peer(map(), atom(), map()) -> ok | {error, term()}.
send_to_peer(NodeInfo, MessageType, Message) ->
Endpoint = extract_endpoint(NodeInfo),
?LOG_DEBUG("[DHT] send_to_peer: NodeInfo=~p, Endpoint=~p, Type=~p",
[maps:get(node_id, NodeInfo, unknown), Endpoint, MessageType]),
do_send_to_peer(Endpoint, MessageType, Message).
do_send_to_peer(undefined, MessageType, _Message) ->
?LOG_WARNING("[DHT] send_to_peer: no endpoint for message type ~p", [MessageType]),
{error, no_endpoint};
do_send_to_peer(Endpoint, MessageType, Message) ->
%% Send directly via peer connector (establishes QUIC connection)
Result = macula_peer_connector:send_message(Endpoint, MessageType, Message),
?LOG_DEBUG("[DHT] send_to_peer result to ~s: ~p", [Endpoint, Result]),
Result.
%% @private
%% @doc Extract endpoint from node info, constructing from address if needed.
extract_endpoint(NodeInfo) ->
extract_endpoint_from_field(maps:get(endpoint, NodeInfo, undefined), NodeInfo).
extract_endpoint_from_field(undefined, NodeInfo) ->
construct_endpoint_from_address(maps:get(address, NodeInfo, undefined));
extract_endpoint_from_field(Endpoint, _NodeInfo) ->
Endpoint.
construct_endpoint_from_address(undefined) ->
undefined;
construct_endpoint_from_address({Host, Port}) when is_integer(Port) ->
HostBin = format_host_to_binary(Host),
PortBin = integer_to_binary(Port),
<<HostBin/binary, ":", PortBin/binary>>;
construct_endpoint_from_address(HostPortStr) when is_binary(HostPortStr) ->
HostPortStr;
construct_endpoint_from_address(HostPortStr) when is_list(HostPortStr) ->
list_to_binary(HostPortStr);
construct_endpoint_from_address(_Other) ->
undefined.
format_host_to_binary({_, _, _, _} = IPv4) ->
list_to_binary(inet:ntoa(IPv4));
format_host_to_binary({_, _, _, _, _, _, _, _} = IPv6) ->
list_to_binary(inet:ntoa(IPv6));
format_host_to_binary(Host) when is_list(Host) ->
list_to_binary(Host);
format_host_to_binary(Host) when is_binary(Host) ->
Host.
%% @doc Query remote peer and wait for response.
%% Used for FIND_NODE and FIND_VALUE operations.
%% Currently uses fire-and-forget delivery. For request/response patterns,
%% use macula_rpc_handler:request/4 which provides NATS-style async RPC
%% with callbacks (available since v0.12.1).
-spec query_peer(map(), atom(), map()) -> {ok, term()} | {error, term()}.
query_peer(NodeInfo, MessageType, Message) ->
send_to_peer(NodeInfo, MessageType, Message).
%% @private
%% @doc Encode reply based on message type.
encode_reply_by_type(find_node, Reply) ->
macula_protocol_encoder:encode(find_node_reply, Reply);
encode_reply_by_type(find_value, Reply) ->
macula_protocol_encoder:encode(find_value_reply, Reply);
encode_reply_by_type(store, Reply) ->
macula_protocol_encoder:encode(store, Reply);
encode_reply_by_type(_UnknownType, _Reply) ->
macula_protocol_encoder:encode(reply, #{error => <<"Unknown DHT message type">>}).