Packages
macula
0.7.13
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/4
]).
%% Exported for testing
-ifdef(TEST).
-export([
deliver_to_local_subscribers/2,
route_to_remote_subscribers/4
]).
-endif.
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Distribute pub/sub message to both local and remote subscribers.
%% Uses DHT routing for remote subscribers (multi-hop Kademlia).
-spec distribute(
LocalSubscribers :: [quicer:stream_handle()],
PubMsg :: map(),
LocalNodeId :: binary(),
Mesh :: pid()
) -> ok.
distribute(LocalSubscribers, PubMsg, LocalNodeId, Mesh) ->
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),
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()) -> ok.
route_to_remote_subscribers(Topic, PubMsg, LocalNodeId, Mesh) ->
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);
{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()) -> ok.
route_to_each_subscriber([], _Topic, _PubMsg, _LocalNodeId, _Mesh) ->
ok;
route_to_each_subscriber([Subscriber | Rest], Topic, PubMsg, LocalNodeId, Mesh) ->
route_to_single_subscriber(Subscriber, Topic, PubMsg, LocalNodeId, Mesh),
route_to_each_subscriber(Rest, Topic, PubMsg, LocalNodeId, Mesh).
%% @private
%% @doc Route message to a single remote subscriber (if node_id present).
-spec route_to_single_subscriber(map(), binary(), map(), binary(), pid()) -> ok.
route_to_single_subscriber(#{<<"node_id">> := DestNodeId}, Topic, PubMsg, LocalNodeId, Mesh) ->
Payload = maps:get(<<"payload">>, PubMsg),
Qos = maps:get(<<"qos">>, PubMsg, 0),
%% Create PUBLISH message for this subscriber
PublishMsg = #{
<<"topic">> => Topic,
<<"payload">> => Payload,
<<"qos">> => Qos,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
},
%% Wrap in pubsub_route envelope and send via DHT
PubSubRouteMsg = macula_pubsub_routing:wrap_publish(LocalNodeId, DestNodeId, PublishMsg, 10),
send_via_dht(DestNodeId, PubSubRouteMsg, Mesh);
route_to_single_subscriber(_Subscriber, _Topic, _PubMsg, _LocalNodeId, _Mesh) ->
io:format("[PubSubRouter] Subscriber missing node_id, skipping~n"),
ok.
%% @private
%% @doc Send pubsub_route message via DHT mesh connection.
-spec send_via_dht(binary(), map(), pid()) -> ok.
send_via_dht(DestNodeId, PubSubRouteMsg, Mesh) ->
case macula_gateway_mesh:get_or_create_connection(Mesh, DestNodeId, undefined) 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.
%% @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.