Packages
macula
0.31.3
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_system/macula_pubsub_dht.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% DHT operations for pub/sub - handles subscription advertisement and discovery.
%%%
%%% Responsibilities:
%%% - Advertise subscriptions in DHT with TTL
%%% - Schedule re-advertisement timers
%%% - Discover remote subscribers via DHT queries
%%% - Route messages to remote subscribers
%%% - Track pending DHT queries
%%%
%%% Extracted from macula_pubsub_handler.erl (Phase 3)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_dht).
-include_lib("kernel/include/logger.hrl").
%% API
-export([advertise_subscription/5, cancel_advertisement/2,
discover_subscribers/6, handle_discovery_response/3,
route_to_subscribers/5]).
-type topic() :: binary().
-type node_id() :: binary().
-type url() :: binary().
-type connection_manager_pid() :: pid().
-type subscription_ref() :: reference().
-type payload() :: binary().
-type qos() :: 0 | 1.
-type advertised_subscriptions() :: #{topic() => #{
sub_ref := reference(),
ttl := pos_integer(),
timer_ref := reference()
}}.
-type pending_queries() :: #{binary() => {topic(), payload(), qos(), map()}}.
-export_type([advertised_subscriptions/0, pending_queries/0]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Advertise a subscription in the DHT.
%% Sends STORE message to DHT and schedules re-advertisement.
%% Returns {ok, SubInfo}.
-spec advertise_subscription(topic(), subscription_ref(), node_id(), url(), connection_manager_pid()) ->
{ok, #{sub_ref := reference(), ttl := pos_integer(), timer_ref := reference()}} | {error, term()}.
advertise_subscription(Topic, SubRef, NodeId, Url, _ConnMgrPid) ->
%% Default TTL for subscriptions is 300 seconds (5 minutes)
TTL = 300,
%% Create DHT key by hashing the topic
TopicKey = crypto:hash(sha256, Topic),
%% Create subscriber value with endpoint info for routing messages back
SubscriberValue = #{
node_id => NodeId,
endpoint => Url,
ttl => TTL
},
%% Store subscription in DHT with propagation to k closest nodes (v0.8.0+)
?LOG_INFO("[~s] Advertising subscription to topic ~s in DHT", [NodeId, Topic]),
store_subscription_in_dht(NodeId, Topic, TopicKey, SubscriberValue),
%% Schedule re-advertisement (TTL - 60 seconds, minimum 10 seconds)
ResubInterval = max(10, TTL - 60) * 1000, % milliseconds
TimerRef = erlang:send_after(ResubInterval, self(), {resubscribe, Topic}),
?LOG_DEBUG("[~s] Scheduled re-subscription for ~s in ~p seconds",
[NodeId, Topic, ResubInterval div 1000]),
%% Return subscription info
SubInfo = #{
sub_ref => SubRef,
ttl => TTL,
timer_ref => TimerRef
},
{ok, SubInfo}.
%% @doc Cancel advertisement for a topic.
%% Cancels the re-advertisement timer.
%% Returns updated advertised_subscriptions map.
-spec cancel_advertisement(topic(), advertised_subscriptions()) -> advertised_subscriptions().
cancel_advertisement(Topic, AdvertisedSubscriptions) ->
SubInfo = maps:get(Topic, AdvertisedSubscriptions, undefined),
do_cancel_advertisement(SubInfo, Topic, AdvertisedSubscriptions).
%% @private Not advertised, return unchanged
do_cancel_advertisement(undefined, _Topic, AdvertisedSubscriptions) ->
AdvertisedSubscriptions;
%% @private Cancel timer and remove from map
do_cancel_advertisement(SubInfo, Topic, AdvertisedSubscriptions) ->
TimerRef = maps:get(timer_ref, SubInfo),
erlang:cancel_timer(TimerRef),
maps:remove(Topic, AdvertisedSubscriptions).
%% @doc Discover remote subscribers for a topic.
%% Checks cache first, queries DHT on cache miss.
%% Returns {cached, Subscribers, Registry} | {query_sent, Pending, MsgId, Registry}.
-spec discover_subscribers(topic(), payload(), qos(), connection_manager_pid(), term(), non_neg_integer()) ->
{cached, list(), term()} | {query_sent, pending_queries(), binary(), term()}.
discover_subscribers(Topic, Payload, Qos, ConnMgrPid, ServiceRegistry, MsgIdCounter) ->
%% Check subscriber cache first
case macula_service_registry:discover_subscribers(ServiceRegistry, Topic) of
{ok, Subscribers, UpdatedRegistry} ->
%% Cache hit - return cached subscribers
?LOG_DEBUG("Cache hit for subscribers to topic: ~s (~p subscribers)",
[Topic, length(Subscribers)]),
{cached, Subscribers, UpdatedRegistry};
{cache_miss, UpdatedRegistry} ->
%% Cache miss - query DHT asynchronously
{MsgId, _NewCounter} = next_message_id(MsgIdCounter),
query_dht_async(Topic, Payload, Qos, MsgId, ConnMgrPid),
%% Track pending query
Pending = #{MsgId => {Topic, Payload, Qos, #{}}},
{query_sent, Pending, MsgId, UpdatedRegistry}
end.
%% @doc Handle DHT discovery response.
%% Routes messages to discovered subscribers.
%% Returns updated pending queries map.
-spec handle_discovery_response(binary(), list(), pending_queries()) ->
{ok, pending_queries()} | {not_found, pending_queries()}.
handle_discovery_response(MsgId, _Subscribers, PendingQueries) ->
QueryInfo = maps:get(MsgId, PendingQueries, undefined),
do_handle_discovery_response(QueryInfo, MsgId, PendingQueries).
%% @private Query not found (already handled or unknown)
do_handle_discovery_response(undefined, _MsgId, PendingQueries) ->
{not_found, PendingQueries};
%% @private Query found - remove from pending
do_handle_discovery_response({_Topic, _Payload, _Qos, _Opts}, MsgId, PendingQueries) ->
{ok, maps:remove(MsgId, PendingQueries)}.
%% @doc Route message to remote subscribers via direct P2P connections (v0.8.0+).
%% Wraps publish in pubsub_route envelope and sends directly to each subscriber.
%% Uses macula_peer_connector for direct QUIC connections to subscriber endpoints.
-spec route_to_subscribers(topic(), payload(), qos(), list(), node_id()) -> ok.
route_to_subscribers(_Topic, _Payload, _Qos, [], _NodeId) ->
ok;
route_to_subscribers(Topic, Payload, Qos, Subscribers, SourceNodeId) ->
?LOG_INFO("[~s] Routing message to ~p remote subscriber(s) for topic: ~s via P2P",
[SourceNodeId, length(Subscribers), Topic]),
%% Route to each subscriber via direct P2P connection
lists:foreach(
fun(Subscriber) ->
%% Extract subscriber node_id and endpoint
NodeId = maps:get(node_id, Subscriber, maps:get(<<"node_id">>, Subscriber, undefined)),
Endpoint = maps:get(endpoint, Subscriber, maps:get(<<"endpoint">>, Subscriber, undefined)),
case {NodeId, Endpoint} of
{undefined, _} ->
?LOG_WARNING("[~s] Subscriber missing node_id, skipping", [SourceNodeId]);
{_, undefined} ->
?LOG_WARNING("[~s] Subscriber missing endpoint, skipping", [SourceNodeId]);
{DestNodeId, DestEndpoint} ->
%% Build PUBLISH message
PublishMsg = #{
<<"topic">> => Topic,
<<"payload">> => Payload,
<<"qos">> => Qos,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
},
%% Wrap in pubsub_route envelope (MaxHops = 10)
PubSubRouteMsg = macula_pubsub_routing:wrap_publish(
SourceNodeId, DestNodeId, PublishMsg, 10
),
%% Try direct send first, fall back to NAT-aware routing (v0.12.0+)
send_to_subscriber(SourceNodeId, DestNodeId, DestEndpoint, Topic, PubSubRouteMsg)
end
end,
Subscribers
),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private Store subscription in DHT
store_subscription_in_dht(NodeId, Topic, TopicKey, SubscriberValue) ->
case whereis(macula_routing_server) of
undefined ->
?LOG_WARNING("[~s] Routing server not running, cannot advertise subscription", [NodeId]);
RoutingServerPid ->
handle_store_result(
macula_routing_server:store(RoutingServerPid, TopicKey, SubscriberValue),
NodeId, Topic
)
end.
%% @private Handle DHT store result
handle_store_result(ok, NodeId, Topic) ->
?LOG_DEBUG("[~s] Successfully stored subscription for ~s in DHT", [NodeId, Topic]);
handle_store_result({error, StoreError}, NodeId, Topic) ->
?LOG_WARNING("[~s] Failed to store subscription ~s: ~p", [NodeId, Topic, StoreError]).
%% @doc Generate next message ID
-spec next_message_id(non_neg_integer()) -> {binary(), non_neg_integer()}.
next_message_id(Counter) ->
macula_utils:next_message_id(Counter).
%% @doc Query DHT for subscribers and route message to them (v0.8.0+).
%% Spawns the DHT lookup + routing to avoid blocking the pubsub handler.
%% The routing server's find_value can do network I/O (QUIC) which may be slow.
-spec query_dht_async(topic(), payload(), qos(), binary(), connection_manager_pid()) -> ok.
query_dht_async(Topic, Payload, Qos, _MsgId, _ConnMgrPid) ->
TopicKey = crypto:hash(sha256, Topic),
spawn(fun() -> do_query_and_route(Topic, Payload, Qos, TopicKey) end),
ok.
%% @private Query routing server and route to discovered subscribers.
do_query_and_route(Topic, Payload, Qos, TopicKey) ->
case whereis(macula_routing_server) of
undefined ->
?LOG_WARNING("Routing server not running, cannot query for subscribers to ~s", [Topic]);
RoutingServerPid ->
handle_dht_query_result(
macula_routing_server:find_value(RoutingServerPid, TopicKey, 20),
Topic, Payload, Qos
)
end.
%% @private Handle DHT query result
handle_dht_query_result({ok, Subscribers}, Topic, Payload, Qos) when is_list(Subscribers), length(Subscribers) > 0 ->
?LOG_INFO("Found ~p subscriber(s) for topic ~s in DHT, routing",
[length(Subscribers), Topic]),
SourceNodeId = get_local_node_id(),
route_to_subscribers(Topic, Payload, Qos, Subscribers, SourceNodeId);
handle_dht_query_result({ok, []}, Topic, _Payload, _Qos) ->
?LOG_WARNING("No subscribers found for topic ~s in DHT", [Topic]);
handle_dht_query_result({ok, SingleSub}, Topic, Payload, Qos) when is_map(SingleSub) ->
?LOG_DEBUG("Found 1 subscriber for topic ~s in DHT, routing message", [Topic]),
SourceNodeId = get_local_node_id(),
route_to_subscribers(Topic, Payload, Qos, [SingleSub], SourceNodeId);
handle_dht_query_result({error, not_found}, Topic, _Payload, _Qos) ->
?LOG_WARNING("No subscribers in DHT for topic ~s (not_found)", [Topic]);
handle_dht_query_result({error, QueryError}, Topic, _Payload, _Qos) ->
?LOG_WARNING("Failed to query DHT for topic ~s: ~p", [Topic, QueryError]).
%% @private Get local node ID for routing.
-spec get_local_node_id() -> node_id().
get_local_node_id() ->
case whereis(macula_gateway) of
undefined ->
crypto:strong_rand_bytes(32);
GatewayPid ->
get_node_id_from_gateway(gen_server:call(GatewayPid, get_node_id, 1000))
end.
%% @private Extract node ID from gateway response
get_node_id_from_gateway({ok, NodeId}) -> NodeId;
get_node_id_from_gateway(_) -> crypto:strong_rand_bytes(32).
%% @doc Send message to subscriber with NAT-aware fallback (v0.12.0+).
%% Tries direct connection first, then falls back to NAT-aware routing.
-spec send_to_subscriber(node_id(), node_id(), binary(), topic(), map()) -> ok.
send_to_subscriber(SourceNodeId, DestNodeId, DestEndpoint, Topic, PubSubRouteMsg) ->
%% Try direct send first (works for public IPs and same-network peers)
case macula_peer_connector:send_message(DestEndpoint, pubsub_route, PubSubRouteMsg) of
ok ->
?LOG_INFO("[~s] Sent pubsub_route to ~s at ~s for topic ~s",
[SourceNodeId, binary:encode_hex(DestNodeId), DestEndpoint, Topic]);
{error, Reason} ->
?LOG_WARNING("[~s] Direct send to ~s (~s) failed: ~p, trying NAT-aware",
[SourceNodeId, DestEndpoint, Topic, Reason]),
%% Fall back to NAT-aware routing (hole punch, relay)
send_to_subscriber_nat_aware(SourceNodeId, DestNodeId, DestEndpoint, Topic, PubSubRouteMsg)
end.
%% @private Use NAT-aware routing to reach subscriber behind NAT.
-spec send_to_subscriber_nat_aware(node_id(), node_id(), binary(), topic(), map()) -> ok.
send_to_subscriber_nat_aware(SourceNodeId, DestNodeId, DestEndpoint, Topic, PubSubRouteMsg) ->
%% Use NAT-aware connector with endpoint hint
Opts = #{endpoint => DestEndpoint},
case macula_peer_connector:send_message_nat_aware(SourceNodeId, DestNodeId, pubsub_route, PubSubRouteMsg, Opts) of
ok ->
?LOG_INFO("[~s] NAT-aware send to subscriber ~s succeeded for topic ~s",
[SourceNodeId, binary:encode_hex(DestNodeId), Topic]);
{error, Reason} ->
?LOG_ERROR("[~s] NAT-aware send to ~s failed: ~p (topic: ~s)",
[SourceNodeId, binary:encode_hex(DestNodeId), Reason, Topic])
end.