Packages

macula

0.15.0
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_pubsub_system macula_pubsub_dht.erl
Raw

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]),
try
case whereis(macula_routing_server) of
undefined ->
?LOG_WARNING("[~s] Routing server not running, cannot advertise subscription", [NodeId]);
RoutingServerPid ->
case macula_routing_server:store(RoutingServerPid, TopicKey, SubscriberValue) of
ok ->
?LOG_DEBUG("[~s] Successfully stored subscription for ~s in DHT", [NodeId, Topic]);
{error, StoreError} ->
?LOG_WARNING("[~s] Failed to store subscription ~s: ~p",
[NodeId, Topic, StoreError])
end
end
catch
_:DhtError ->
?LOG_WARNING("[~s] Failed to advertise subscription ~s in DHT: ~p (continuing)",
[NodeId, Topic, DhtError])
end,
%% 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
%%%===================================================================
%% @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) ->
%% Create DHT key from topic
TopicKey = crypto:hash(sha256, Topic),
?LOG_DEBUG("Querying DHT for remote subscribers to topic: ~s", [Topic]),
%% Query local routing server directly (v0.8.0+)
try
case whereis(macula_routing_server) of
undefined ->
?LOG_WARNING("Routing server not running, cannot query for subscribers to ~s", [Topic]);
RoutingServerPid ->
%% K=20 is standard Kademlia replication factor
case macula_routing_server:find_value(RoutingServerPid, TopicKey, 20) of
{ok, Subscribers} when is_list(Subscribers), length(Subscribers) > 0 ->
?LOG_INFO("Found ~p subscriber(s) for topic ~s in DHT, routing message",
[length(Subscribers), Topic]),
%% Get our node ID for routing
SourceNodeId = get_local_node_id(),
%% Route message to discovered subscribers
route_to_subscribers(Topic, Payload, Qos, Subscribers, SourceNodeId);
{ok, []} ->
?LOG_DEBUG("No subscribers found for topic ~s in DHT", [Topic]);
{ok, SingleSub} when is_map(SingleSub) ->
%% Single subscriber returned as map (not list)
?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);
{error, not_found} ->
?LOG_DEBUG("No subscribers found for topic ~s in DHT", [Topic]);
{error, QueryError} ->
?LOG_WARNING("Failed to query DHT for topic ~s: ~p",
[Topic, QueryError])
end
end
catch
_:Error:Stack ->
?LOG_WARNING("Failed to query DHT for topic ~s: ~p~nStack: ~p",
[Topic, Error, Stack])
end,
ok.
%% @private Get local node ID for routing.
-spec get_local_node_id() -> node_id().
get_local_node_id() ->
case whereis(macula_gateway) of
undefined ->
%% Fallback: generate a temporary node ID
crypto:strong_rand_bytes(32);
GatewayPid ->
try
{ok, NodeId} = gen_server:call(GatewayPid, get_node_id, 1000),
NodeId
catch
_:_ ->
crypto:strong_rand_bytes(32)
end
end.
%% @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.