Packages
macula
0.8.14
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]),
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) ->
case maps:get(Topic, AdvertisedSubscriptions, undefined) of
undefined ->
%% Not advertised, return unchanged
AdvertisedSubscriptions;
SubInfo ->
%% Cancel timer and remove from map
TimerRef = maps:get(timer_ref, SubInfo),
erlang:cancel_timer(TimerRef),
maps:remove(Topic, AdvertisedSubscriptions)
end.
%% @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) ->
case maps:get(MsgId, PendingQueries, undefined) of
undefined ->
%% Query not found (already handled or unknown)
{not_found, PendingQueries};
{_Topic, _Payload, _Qos, _Opts} ->
%% Query found - remove from pending
%% Note: Actual routing happens in the caller
{ok, maps:remove(MsgId, PendingQueries)}
end.
%% @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
),
%% Send directly to subscriber via peer connector
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_ERROR("[~s] Failed to send pubsub_route to ~s: ~p",
[SourceNodeId, DestEndpoint, Reason])
end
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 synchronously (v0.8.0+).
%% Directly queries local routing server instead of sending message to connection manager.
-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+)
%% This is actually synchronous now, but we keep the function name for compatibility
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_DEBUG("Found ~p subscriber(s) for topic ~s in DHT",
[length(Subscribers), Topic]);
{ok, []} ->
?LOG_DEBUG("No subscribers found for topic ~s in DHT", [Topic]);
{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.