Packages
macula
0.20.22
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+).
%% Directly queries local routing server and routes message to discovered subscribers.
%% NOTE: Despite the "async" name (kept for compatibility), this now performs
%% synchronous DHT lookup AND routing in one call.
-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),
?LOG_DEBUG("Querying DHT for remote subscribers to topic: ~s", [Topic]),
query_routing_server_for_subscribers(Topic, Payload, Qos, TopicKey),
ok.
%% @private Query routing server for subscribers
query_routing_server_for_subscribers(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 message",
[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_DEBUG("No subscribers found for topic ~s in DHT", [Topic]);
handle_dht_query_result({ok, SingleSub}, Topic, Payload, Qos) when is_map(SingleSub) ->
?LOG_INFO("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_DEBUG("No subscribers found for topic ~s in DHT", [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_DEBUG("[~s] Sent pubsub_route directly to subscriber ~s at ~s for topic ~s",
[SourceNodeId, binary:encode_hex(DestNodeId), DestEndpoint, Topic]);
{error, Reason} ->
?LOG_DEBUG("[~s] Direct send to ~s failed (~p), trying NAT-aware routing",
[SourceNodeId, DestEndpoint, 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.