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_pubsub_routing.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Pub/Sub routing for multi-hop DHT-routed pub/sub.
%%% Handles wrapping, unwrapping, and routing of PUBLISH messages through
%%% the Kademlia DHT mesh.
%%%
%%% Pattern: Clone of macula_rpc_routing for pub/sub messages
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_routing).
%% API
-export([
wrap_publish/4,
route_or_deliver/3,
should_deliver_locally/2
]).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Wrap a PUBLISH message in pubsub_route envelope for DHT routing.
-spec wrap_publish(binary(), binary(), macula_protocol_types:publish_msg(), pos_integer()) ->
macula_protocol_types:pubsub_route_msg().
wrap_publish(SourceNodeId, DestinationNodeId, PublishMsg, MaxHops)
when is_binary(SourceNodeId), is_binary(DestinationNodeId) ->
%% Extract topic from publish message
Topic = maps:get(<<"topic">>, PublishMsg),
#{
<<"destination_node_id">> => DestinationNodeId,
<<"source_node_id">> => SourceNodeId,
<<"hop_count">> => 0,
<<"max_hops">> => MaxHops,
<<"topic">> => Topic,
<<"payload">> => PublishMsg
}.
%% @doc Determine if this node should deliver the message locally or forward it.
-spec should_deliver_locally(binary(), macula_protocol_types:pubsub_route_msg()) -> boolean().
should_deliver_locally(LocalNodeId, PubSubRouteMsg) ->
%% MessagePack decoder returns binary keys
#{<<"destination_node_id">> := DestNodeId} = PubSubRouteMsg,
LocalNodeId =:= DestNodeId.
%% @doc Route a pubsub_route message: either deliver locally or forward to next hop.
%% Returns one of:
%% {deliver, Topic, PublishMsg} - Message is for this node
%% {forward, NextHopNodeInfo, UpdatedPubSubRouteMsg} - Forward to next hop
%% {error, Reason} - Cannot route (TTL exceeded, no route, etc.)
-spec route_or_deliver(binary(), macula_protocol_types:pubsub_route_msg(), pid()) ->
{deliver, binary(), map()} |
{forward, macula_routing_bucket:node_info(), macula_protocol_types:pubsub_route_msg()} |
{error, term()}.
route_or_deliver(LocalNodeId, PubSubRouteMsg, RoutingServerPid) ->
%% MessagePack decoder returns binary keys, not atoms
#{<<"destination_node_id">> := DestNodeId,
<<"source_node_id">> := SourceNodeId,
<<"hop_count">> := HopCount,
<<"max_hops">> := MaxHops,
<<"topic">> := Topic,
<<"payload">> := Payload} = PubSubRouteMsg,
%% Check TTL and route
check_hop_count(HopCount >= MaxHops, LocalNodeId, DestNodeId, SourceNodeId,
HopCount, MaxHops, Topic, Payload, PubSubRouteMsg, RoutingServerPid).
%% @doc Check if hop count exceeded.
check_hop_count(true, _LocalNodeId, DestNodeId, SourceNodeId, HopCount, MaxHops, _Topic, _Payload, _PubSubRouteMsg, _RoutingServerPid) ->
?LOG_WARNING("Pub/Sub route exceeded max hops (~p >= ~p) from ~p to ~p",
[HopCount, MaxHops, SourceNodeId, DestNodeId]),
{error, max_hops_exceeded};
check_hop_count(false, LocalNodeId, DestNodeId, SourceNodeId, HopCount, _MaxHops, Topic, Payload, PubSubRouteMsg, RoutingServerPid) ->
check_local_delivery(LocalNodeId =:= DestNodeId, DestNodeId, SourceNodeId,
HopCount, Topic, Payload, PubSubRouteMsg, RoutingServerPid).
%% @doc Check if this is local delivery.
check_local_delivery(true, _DestNodeId, SourceNodeId, HopCount, Topic, Payload, _PubSubRouteMsg, _RoutingServerPid) ->
%% Deliver locally
?LOG_DEBUG("Pub/Sub route delivering locally: topic ~p from ~p (hops: ~p)",
[Topic, SourceNodeId, HopCount]),
{deliver, Topic, Payload};
check_local_delivery(false, DestNodeId, _SourceNodeId, _HopCount, _Topic, _Payload, PubSubRouteMsg, RoutingServerPid) ->
%% Forward to next hop
forward_to_next_hop(DestNodeId, PubSubRouteMsg, RoutingServerPid).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Find next hop closer to destination and prepare forwarding.
-spec forward_to_next_hop(binary(), macula_protocol_types:pubsub_route_msg(), pid()) ->
{forward, macula_routing_bucket:node_info(), macula_protocol_types:pubsub_route_msg()} |
{error, term()}.
forward_to_next_hop(DestNodeId, PubSubRouteMsg, RoutingServerPid) ->
%% Query routing table for closest node to destination
%% K=3 gives us some redundancy if first hop fails
case macula_routing_server:find_closest(RoutingServerPid, DestNodeId, 3) of
[] ->
?LOG_ERROR("No route to destination: ~p", [DestNodeId]),
{error, no_route};
[NextHop | _] ->
%% Increment hop count (MessagePack uses binary keys)
#{<<"hop_count">> := HopCount} = PubSubRouteMsg,
UpdatedMsg = PubSubRouteMsg#{<<"hop_count">> => HopCount + 1},
?LOG_DEBUG("Forwarding Pub/Sub route to next hop ~p (destination: ~p, hops: ~p)",
[maps:get(node_id, NextHop), DestNodeId, HopCount + 1]),
{forward, NextHop, UpdatedMsg}
end.