Packages
macula
0.7.25
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_gateway_pubsub_router.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Macula Gateway Pub/Sub Router - DHT-Routed Message Distribution
%%%
%%% Handles distribution of pub/sub messages to both local and remote
%%% subscribers using multi-hop Kademlia DHT routing (v0.7.8+).
%%%
%%% Responsibilities:
%%% - Deliver messages to local subscribers via QUIC streams
%%% - Query DHT for remote subscribers
%%% - Route messages via DHT multi-hop (pubsub_route protocol)
%%% - Wrap PUBLISH messages in pubsub_route envelopes
%%%
%%% Extracted from macula_gateway.erl (v0.7.9) for better separation of concerns.
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_pubsub_router).
%% API
-export([
distribute/5
]).
%% Exported for testing
-ifdef(TEST).
-export([
deliver_to_local_subscribers/2,
route_to_remote_subscribers/5,
parse_endpoint/1,
parse_endpoint_host/2
]).
-endif.
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Distribute pub/sub message to both local and remote subscribers.
%% Uses DHT routing for remote subscribers (multi-hop Kademlia).
%% For connected clients, uses existing bidirectional streams instead of mesh connections.
-spec distribute(
LocalSubscribers :: [quicer:stream_handle()],
PubMsg :: map(),
LocalNodeId :: binary(),
Mesh :: pid(),
Clients :: pid()
) -> ok.
distribute(LocalSubscribers, PubMsg, LocalNodeId, Mesh, Clients) ->
Topic = maps:get(<<"topic">>, PubMsg),
%% Deliver to local and remote subscribers
deliver_to_local_subscribers(LocalSubscribers, PubMsg),
route_to_remote_subscribers(Topic, PubMsg, LocalNodeId, Mesh, Clients),
io:format("[PubSubRouter] Finished distributing message to ~p local + DHT remote subscribers~n",
[length(LocalSubscribers)]),
ok.
%%%===================================================================
%%% Internal Functions - Local Delivery
%%%===================================================================
%% @private
%% @doc Deliver message to local subscribers via QUIC streams.
-spec deliver_to_local_subscribers([quicer:stream_handle()], map()) -> ok.
deliver_to_local_subscribers([], _PubMsg) ->
ok;
deliver_to_local_subscribers(Subscribers, PubMsg) ->
PubBinary = macula_protocol_encoder:encode(publish, PubMsg),
io:format("[PubSubRouter] Found ~p local subscribers~n", [length(Subscribers)]),
lists:foreach(fun(Stream) -> send_to_local_stream(Stream, PubBinary) end, Subscribers),
ok.
%% @private
%% @doc Send message to a single local subscriber stream.
-spec send_to_local_stream(quicer:stream_handle(), binary()) -> ok.
send_to_local_stream(Stream, Binary) ->
case macula_quic:send(Stream, Binary) of
ok ->
io:format("[PubSubRouter] Successfully sent to local stream ~p~n", [Stream]),
ok;
{error, Reason} ->
io:format("[PubSubRouter] Failed to send to local stream ~p: ~p~n", [Stream, Reason]),
ok
end.
%%%===================================================================
%%% Internal Functions - DHT Remote Routing
%%%===================================================================
%% @private
%% @doc Route message to remote subscribers via DHT multi-hop routing.
-spec route_to_remote_subscribers(binary(), map(), binary(), pid(), pid()) -> ok.
route_to_remote_subscribers(Topic, PubMsg, LocalNodeId, Mesh, Clients) ->
TopicKey = crypto:hash(sha256, Topic),
case macula_gateway_dht:lookup_value(TopicKey) of
{ok, RemoteSubscribers} ->
io:format("[PubSubRouter] Found ~p remote subscriber(s) in DHT~n", [length(RemoteSubscribers)]),
route_to_each_subscriber(RemoteSubscribers, Topic, PubMsg, LocalNodeId, Mesh, Clients);
{error, not_found} ->
io:format("[PubSubRouter] No remote subscribers found in DHT~n"),
ok
end.
%% @private
%% @doc Route message to each remote subscriber via DHT.
-spec route_to_each_subscriber([map()], binary(), map(), binary(), pid(), pid()) -> ok.
route_to_each_subscriber([], _Topic, _PubMsg, _LocalNodeId, _Mesh, _Clients) ->
ok;
route_to_each_subscriber([Subscriber | Rest], Topic, PubMsg, LocalNodeId, Mesh, Clients) ->
route_to_single_subscriber(Subscriber, Topic, PubMsg, LocalNodeId, Mesh, Clients),
route_to_each_subscriber(Rest, Topic, PubMsg, LocalNodeId, Mesh, Clients).
%% @private
%% @doc Route message to a single remote subscriber (if node_id present).
%% Checks if subscriber is a connected client first - if so, uses existing bidirectional stream.
%% Otherwise creates mesh connection.
-spec route_to_single_subscriber(map(), binary(), map(), binary(), pid(), pid()) -> ok.
route_to_single_subscriber(#{<<"node_id">> := DestNodeId} = Subscriber, Topic, PubMsg, LocalNodeId, Mesh, Clients) ->
Payload = maps:get(<<"payload">>, PubMsg),
Qos = maps:get(<<"qos">>, PubMsg, 0),
%% DEBUG: Log the node_id from DHT (should be raw 32-byte binary)
io:format("[PubSubRouter DEBUG] DestNodeId size: ~p bytes, hex: ~s~n",
[byte_size(DestNodeId), binary:encode_hex(DestNodeId)]),
%% Create PUBLISH message for this subscriber
PublishMsg = #{
<<"topic">> => Topic,
<<"payload">> => Payload,
<<"qos">> => Qos,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
},
%% Check if this is a connected client (has existing bidirectional stream)
case macula_gateway_clients:get_client_stream(Clients, DestNodeId) of
{ok, Stream} ->
%% Subscriber is a connected client - send directly via their stream
io:format("[PubSubRouter] ✅ Routing to connected client ~s via existing stream~n",
[binary:encode_hex(DestNodeId)]),
send_to_client_stream(Stream, PublishMsg);
not_found ->
%% Not a connected client - must be remote gateway, use mesh connection
io:format("[PubSubRouter] ❌ Client stream NOT FOUND for ~s, routing via mesh~n",
[binary:encode_hex(DestNodeId)]),
%% DEBUG: List all stored client streams to compare
io:format("[PubSubRouter DEBUG] Getting all client stream keys...~n"),
AllKeys = macula_gateway_clients:get_all_node_ids(Clients),
io:format("[PubSubRouter DEBUG] Stored client stream node_ids:~n"),
lists:foreach(fun(KeyNodeId) ->
io:format(" - ~s (size: ~p bytes)~n", [binary:encode_hex(KeyNodeId), byte_size(KeyNodeId)])
end, AllKeys),
io:format("[PubSubRouter DEBUG] Looking for: ~s (size: ~p bytes)~n",
[binary:encode_hex(DestNodeId), byte_size(DestNodeId)]),
EndpointUrl = maps:get(<<"endpoint">>, Subscriber, undefined),
PubSubRouteMsg = macula_pubsub_routing:wrap_publish(LocalNodeId, DestNodeId, PublishMsg, 10),
send_via_dht(DestNodeId, EndpointUrl, PubSubRouteMsg, Mesh)
end;
route_to_single_subscriber(Subscriber, _Topic, _PubMsg, _LocalNodeId, _Mesh, _Clients) ->
io:format("[PubSubRouter] Subscriber missing node_id, skipping: ~p~n", [Subscriber]),
ok.
%% @private
%% @doc Send pubsub_route message via DHT mesh connection.
%% Parses endpoint URL to extract address and port.
-spec send_via_dht(binary(), binary(), map(), pid()) -> ok.
send_via_dht(DestNodeId, EndpointUrl, PubSubRouteMsg, Mesh) ->
case parse_endpoint(EndpointUrl) of
{ok, {Address, Port}} ->
case macula_gateway_mesh:get_or_create_connection(Mesh, DestNodeId, {Address, Port}) of
{ok, Stream} ->
send_route_message(Stream, PubSubRouteMsg, DestNodeId);
{error, Reason} ->
io:format("[PubSubRouter] Failed to get connection for ~s: ~p~n",
[binary:encode_hex(DestNodeId), Reason]),
ok
end;
{error, ParseReason} ->
io:format("[PubSubRouter] Failed to parse endpoint ~s: ~p~n",
[EndpointUrl, ParseReason]),
ok
end.
%% @private
%% @doc Send PUBLISH message directly to connected client stream.
-spec send_to_client_stream(quicer:stream_handle(), map()) -> ok.
send_to_client_stream(Stream, PublishMsg) ->
PubBinary = macula_protocol_encoder:encode(publish, PublishMsg),
case macula_quic:send(Stream, PubBinary) of
ok ->
io:format("[PubSubRouter] Sent PUBLISH to connected client stream~n"),
ok;
{error, Reason} ->
io:format("[PubSubRouter] Failed to send to client stream: ~p~n", [Reason]),
ok
end.
%% @private
%% @doc Send encoded pubsub_route message to stream.
-spec send_route_message(quicer:stream_handle(), map(), binary()) -> ok.
send_route_message(Stream, PubSubRouteMsg, DestNodeId) ->
RouteMsg = macula_protocol_encoder:encode(pubsub_route, PubSubRouteMsg),
case macula_quic:send(Stream, RouteMsg) of
ok ->
io:format("[PubSubRouter] Sent pubsub_route to ~s via DHT~n",
[binary:encode_hex(DestNodeId)]),
ok;
{error, Reason} ->
io:format("[PubSubRouter] Failed to send pubsub_route: ~p~n", [Reason]),
ok
end.
%%%===================================================================
%%% Node ID Helpers (Removed - not needed)
%%%===================================================================
%% Per Kademlia spec: node IDs are always 256-bit (32-byte) raw binary.
%% Hex encoding is ONLY for display/logging, never for internal storage.
%% DHT stores and returns raw binary node IDs.
%%%===================================================================
%%% Endpoint Parsing Helper
%%%===================================================================
%% @private
%% @doc Parse endpoint URL (e.g., "https://192.168.1.100:4433") to {Address, Port}.
%% Returns {ok, {Address, Port}} or {error, Reason}.
-spec parse_endpoint(binary()) -> {ok, {inet:ip_address() | list(), inet:port_number()}} | {error, term()}.
parse_endpoint(EndpointUrl) when is_binary(EndpointUrl) ->
case uri_string:parse(EndpointUrl) of
#{host := Host, port := Port} when is_integer(Port) ->
parse_endpoint_host(Host, Port);
#{host := Host} ->
parse_endpoint_host(Host, 4433); % Default QUIC port
_Other ->
{error, {invalid_endpoint_format, EndpointUrl}}
end;
parse_endpoint(_Other) ->
{error, invalid_endpoint_type}.
%% @private
%% Parse hostname to IP address tuple.
parse_endpoint_host(Host, Port) when is_list(Host) ->
case inet:parse_address(Host) of
{ok, IpTuple} ->
{ok, {IpTuple, Port}};
{error, _} ->
%% Not an IP address, treat as hostname
{ok, {Host, Port}}
end;
parse_endpoint_host(Host, Port) when is_binary(Host) ->
parse_endpoint_host(binary_to_list(Host), Port).