Packages

macula

0.7.19
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
macula src macula_gateway_pubsub_router.erl
Raw

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">> := DestNodeIdRaw} = Subscriber, Topic, PubMsg, LocalNodeId, Mesh, Clients) ->
Payload = maps:get(<<"payload">>, PubMsg),
Qos = maps:get(<<"qos">>, PubMsg, 0),
%% Convert node_id to binary (DHT may return as binary or list)
DestNodeId = ensure_binary_node_id(DestNodeIdRaw),
%% 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] Routing to remote gateway ~s via mesh~n",
[binary:encode_hex(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
%%%===================================================================
%% @private
%% @doc Ensure node_id is in binary format (not hex string).
%% DHT may return node_id as binary or as hex string.
-spec ensure_binary_node_id(binary() | list()) -> binary().
ensure_binary_node_id(NodeId) when is_binary(NodeId) ->
%% Check if it's a hex string by trying to decode
try binary:decode_hex(NodeId) of
Decoded -> Decoded
catch
error:badarg -> NodeId % Not hex, already raw binary
end;
ensure_binary_node_id(NodeId) when is_list(NodeId) ->
%% Convert list to binary first, then decode
binary:decode_hex(list_to_binary(NodeId)).
%%%===================================================================
%%% 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).