Packages

macula

0.7.27
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_dht.erl
Raw

src/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
},
%% Send STORE message to DHT via connection manager
?LOG_INFO("[~s] Advertising subscription to topic ~s in DHT", [NodeId, Topic]),
try
StoreMsg = macula_routing_protocol:encode_store(TopicKey, SubscriberValue),
case macula_connection:send_message(ConnMgrPid, store, StoreMsg) of
ok ->
?LOG_DEBUG("[~s] Successfully stored subscription for ~s in DHT", [NodeId, Topic]);
{error, SendError} ->
?LOG_WARNING("[~s] Failed to send STORE for subscription ~s: ~p",
[NodeId, Topic, SendError])
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 DHT routing (v0.7.8+).
%% Wraps publish in pubsub_route envelope and sends to each subscriber node.
-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 DHT",
[SourceNodeId, length(Subscribers), Topic]),
%% Get connection manager PID to send routed messages
ConnMgrPid = gproc:lookup_local_name(macula_connection),
%% Route to each subscriber via pubsub_route envelope
lists:foreach(
fun(Subscriber) ->
%% Extract subscriber node_id (not endpoint - we route via DHT)
case maps:get(node_id, Subscriber, maps:get(<<"node_id">>, Subscriber, undefined)) of
undefined ->
?LOG_WARNING("[~s] Subscriber missing node_id, skipping", [SourceNodeId]);
DestNodeId ->
%% 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 via connection manager
case macula_connection:send_message(ConnMgrPid, pubsub_route, PubSubRouteMsg) of
ok ->
?LOG_DEBUG("[~s] Sent pubsub_route to subscriber ~s for topic ~s",
[SourceNodeId, binary:encode_hex(DestNodeId), Topic]);
{error, Reason} ->
?LOG_ERROR("[~s] Failed to send pubsub_route: ~p", [SourceNodeId, 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 asynchronously
-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 (MsgId: ~s)", [Topic, MsgId]),
%% Send FIND_VALUE query to DHT asynchronously (spawn to avoid blocking)
spawn(fun() ->
try
FindValueMsg = macula_routing_protocol:encode_find_value(TopicKey),
case macula_connection:send_message(ConnMgrPid, find_value, FindValueMsg) of
ok ->
?LOG_DEBUG("Sent FIND_VALUE for topic ~s", [Topic]);
{error, SendError} ->
?LOG_WARNING("Failed to send FIND_VALUE for topic ~s: ~p",
[Topic, SendError])
end
catch
_:QueryError:Stack ->
?LOG_WARNING("Failed to query DHT for topic ~s: ~p~nStack: ~p",
[Topic, QueryError, Stack])
end
end),
ok.