Packages
macula
0.20.13
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_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,
forward_publish_to_bootstrap/1,
send_to_peer/3,
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
do_handle_find_value(Key) when is_binary(Key) ->
%% Use routing server's find_value API with k=20
case macula_routing_server:find_value(macula_routing_server, Key, 20) of
{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});
{ok, []} ->
?LOG_DEBUG("[Gateway DHT] FIND_VALUE: no values found"),
macula_routing_protocol:encode_find_value_reply({nodes, []});
{error, not_found} ->
?LOG_DEBUG("[Gateway DHT] FIND_VALUE: not found"),
macula_routing_protocol:encode_find_value_reply({nodes, []});
{error, Reason} ->
?LOG_WARNING("[Gateway DHT] FIND_VALUE error: ~p", [Reason]),
#{type => error, reason => Reason}
end;
do_handle_find_value(Key) ->
?LOG_WARNING("[Gateway DHT] FIND_VALUE with non-binary key: ~p", [Key]),
#{type => error, reason => invalid_key}.
%% @doc Handle DHT FIND_NODE message.
%% Forwards to routing server and sends encoded reply over stream.
%% Crashes on routing server or encoding failures - exposes DHT/protocol bugs.
-spec handle_find_node(reference(), map()) -> ok.
handle_find_node(Stream, FindNodeMsg) ->
%% Forward to routing server (let it crash on errors)
Reply = macula_routing_server:handle_message(macula_routing_server, FindNodeMsg),
%% Send reply back over stream (let it crash on errors)
ReplyBinary = macula_protocol_encoder:encode(find_node_reply, Reply),
macula_quic:send(Stream, ReplyBinary),
ok.
%% @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 a value from the DHT by key.
%% Synchronous lookup from local DHT storage.
%% Subscriptions are replicated via DHT propagation to k closest nodes,
%% so local lookup returns subscribers from the replicated DHT data.
%% Returns list of subscribers for the given key.
-spec lookup_value(binary()) -> {ok, list()} | {error, not_found}.
lookup_value(Key) ->
?LOG_DEBUG("lookup_value called with key hash: ~p", [Key]),
do_lookup_value(whereis(macula_routing_server), Key).
do_lookup_value(undefined, _Key) ->
?LOG_ERROR("routing_server not found!"),
{error, not_found};
do_lookup_value(RoutingServerPid, Key) ->
%% K=20 is the standard Kademlia replication factor
?LOG_DEBUG("Calling routing_server:find_value..."),
Result = macula_routing_server:find_value(RoutingServerPid, Key, 20),
?LOG_DEBUG("find_value result: ~p", [Result]),
normalize_find_value_result(Result).
normalize_find_value_result({ok, []}) ->
{error, not_found};
normalize_find_value_result({ok, Value}) when is_list(Value) ->
{ok, Value};
normalize_find_value_result({ok, Value}) ->
%% Single value, wrap in list
{ok, [Value]};
normalize_find_value_result({error, Reason}) ->
{error, Reason}.
%% @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 DHT message to remote peer (fire-and-forget).
%% Used for STORE operations that don't need a response.
-spec send_to_peer(map(), atom(), map()) -> ok | {error, term()}.
send_to_peer(NodeInfo, MessageType, Message) ->
Endpoint = extract_endpoint(NodeInfo),
?LOG_INFO("[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_INFO("[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">>}).