Packages

macula

0.8.20
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_handler.erl
Raw

src/macula_pubsub_system/macula_pubsub_handler.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% PubSub handler GenServer - facade that orchestrates pub/sub operations.
%%%
%%% This module acts as a facade/coordinator, delegating business logic to:
%%% - macula_pubsub_subscription: Subscription storage, pattern matching, callbacks
%%% - macula_pubsub_dht: DHT advertisement, discovery, routing
%%% - macula_pubsub_qos: QoS 1 tracking and retry logic
%%%
%%% Responsibilities:
%%% - API facade for subscribe/unsubscribe/publish operations
%%% - Message routing coordination between specialized modules
%%% - GenServer lifecycle management
%%% - State management (delegates actual operations to modules)
%%%
%%% Extracted from macula_connection.erl (Phase 4)
%%% Refactored using TDD to extract god module (Phase 5)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_handler).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
%% API
-export([start_link/1, subscribe/3, unsubscribe/2, publish/4, handle_incoming_publish/2]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
opts :: map(),
node_id :: binary(),
url :: binary(),
realm :: binary(),
%% Connection manager PID (looked up via gproc)
connection_manager_pid :: pid() | undefined,
%% Subscriptions: #{SubscriptionRef => {Topic, Callback}}
subscriptions = #{} :: #{reference() => {binary(), fun((map()) -> ok)}},
%% Advertised subscriptions with re-advertisement timers
%% #{Topic => #{sub_ref, ttl, timer_ref}}
advertised_subscriptions = #{} :: #{binary() => #{
sub_ref := reference(),
ttl := pos_integer(),
timer_ref := reference()
}},
%% Pending publish acknowledgments for QoS 1
%% #{MsgId => {Topic, Payload, QoS, RetryCount, TimerRef}}
pending_pubacks = #{} :: #{binary() => {binary(), binary(), integer(), integer(), reference()}},
%% Pending subscriber queries for DHT discovery
%% #{MsgId => {Topic, Payload, QoS, Opts}}
pending_subscriber_queries = #{} :: #{binary() => {binary(), binary(), integer(), map()}},
%% Message ID counter
msg_id_counter = 0 :: non_neg_integer(),
%% Topic pattern matching configuration
topic_separator = <<".">> :: binary(),
topic_wildcard_single = <<"*">> :: binary(),
topic_wildcard_multi = <<"**">> :: binary(),
%% Service registry for DHT operations
service_registry :: macula_service_registry:registry()
}).
%%%===================================================================
%%% API
%%%===================================================================
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
%% No local registration - allows multiple connection instances
gen_server:start_link(?MODULE, Opts, []).
-spec subscribe(pid(), binary() | list() | atom(), fun((map()) -> ok)) -> {ok, reference()} | {error, term()}.
subscribe(Pid, Topic, Callback) ->
gen_server:call(Pid, {subscribe, Topic, Callback}, 5000).
-spec unsubscribe(pid(), reference()) -> ok | {error, term()}.
unsubscribe(Pid, SubRef) ->
gen_server:call(Pid, {unsubscribe, SubRef}, 5000).
-spec publish(pid(), binary() | list() | atom(), term(), map()) -> ok | {error, term()}.
publish(Pid, Topic, Data, Opts) ->
gen_server:call(Pid, {publish, Topic, Data, Opts}, 5000).
-spec handle_incoming_publish(pid(), map()) -> ok.
handle_incoming_publish(Pid, Msg) ->
gen_server:cast(Pid, {incoming_publish, Msg}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
?LOG_INFO("PubSub handler starting"),
%% Extract required options
NodeId = maps:get(node_id, Opts, generate_node_id()),
Url = maps:get(url, Opts, <<"unknown">>),
Realm = maps:get(realm, Opts, <<"default">>),
%% Register in gproc for incoming publish message routing
true = gproc:reg({n, l, {pubsub_handler, Realm}}),
?LOG_INFO("PubSub handler registered in gproc for realm ~s", [Realm]),
%% connection_manager_pid will be set via cast message after init
ConnMgrPid = undefined,
%% Initialize service registry for DHT operations
Registry = macula_service_registry:new(),
%% Topic matching configuration (customizable)
TopicSeparator = maps:get(topic_separator, Opts, <<".">>),
WildcardSingle = maps:get(topic_wildcard_single, Opts, <<"*">>),
WildcardMulti = maps:get(topic_wildcard_multi, Opts, <<"**">>),
{ok, #state{
opts = Opts,
node_id = NodeId,
url = Url,
realm = Realm,
connection_manager_pid = ConnMgrPid,
service_registry = Registry,
topic_separator = TopicSeparator,
topic_wildcard_single = WildcardSingle,
topic_wildcard_multi = WildcardMulti
}}.
%%%===================================================================
%%% Subscribe/Unsubscribe
%%%===================================================================
handle_call({subscribe, Topic, Callback}, _From, State) ->
%% Generate subscription reference
SubRef = make_ref(),
BinaryTopic = ensure_binary(Topic),
%% Send subscribe message via connection manager
SubscribeMsg = #{
topics => [BinaryTopic],
qos => 0
},
case macula_connection:send_message(State#state.connection_manager_pid, subscribe, SubscribeMsg) of
ok ->
%% Store subscription locally (delegate to subscription module)
{ok, UpdatedSubscriptions, SubRef} = macula_pubsub_subscription:add_subscription(
BinaryTopic, Callback, State#state.subscriptions, SubRef
),
State2 = State#state{subscriptions = UpdatedSubscriptions},
%% Advertise subscription in DHT (delegate to DHT module)
{ok, SubInfo} = macula_pubsub_dht:advertise_subscription(
BinaryTopic, SubRef, State#state.node_id, State#state.url,
State#state.connection_manager_pid
),
AdvertisedSubscriptions = State2#state.advertised_subscriptions,
State3 = State2#state{
advertised_subscriptions = AdvertisedSubscriptions#{BinaryTopic => SubInfo}
},
{reply, {ok, SubRef}, State3};
{error, Reason} ->
{reply, {error, Reason}, State}
end;
handle_call({unsubscribe, SubRef}, _From, State) ->
%% Remove subscription (delegate to subscription module)
case macula_pubsub_subscription:remove_subscription(SubRef, State#state.subscriptions) of
{error, not_found} ->
{reply, {error, not_subscribed}, State};
{ok, UpdatedSubscriptions, Topic} ->
%% Send unsubscribe message via connection manager
UnsubscribeMsg = #{
topics => [Topic]
},
case macula_connection:send_message(State#state.connection_manager_pid, unsubscribe, UnsubscribeMsg) of
ok ->
%% Cancel DHT advertisement (delegate to DHT module)
UpdatedAdvertised = macula_pubsub_dht:cancel_advertisement(
Topic, State#state.advertised_subscriptions
),
State2 = State#state{
subscriptions = UpdatedSubscriptions,
advertised_subscriptions = UpdatedAdvertised
},
{reply, ok, State2};
{error, Reason} ->
{reply, {error, Reason}, State}
end
end;
%%%===================================================================
%%% Publish
%%%===================================================================
handle_call({publish, Topic, Data, Opts}, _From, State) ->
%% Check if we have a connection manager and are connected
case State#state.connection_manager_pid of
undefined ->
%% No connection manager - cannot publish
{reply, {error, not_connected}, State};
ConnMgrPid ->
%% Check connection status
case macula_connection:get_status(ConnMgrPid) of
connected ->
%% Build publish message
Qos = maps:get(qos, Opts, 0),
Retain = maps:get(retain, Opts, false),
{MsgId, State2} = next_message_id(State),
BinaryTopic = ensure_binary(Topic),
%% Encode data to binary if it's a map or list
Payload = encode_payload(Data),
PublishMsg = #{
topic => BinaryTopic,
payload => Payload,
qos => Qos,
retain => Retain,
message_id => MsgId
},
%% Send publish message asynchronously using cast to self
%% This avoids blocking the caller on QUIC send operations
gen_server:cast(self(), {do_publish, PublishMsg, Qos, BinaryTopic, Payload, Opts, MsgId}),
{reply, ok, State2};
_Status ->
%% Not connected (connecting, disconnected, or error)
{reply, {error, not_connected}, State}
end
end;
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%%%===================================================================
%%% Async Publish Operations
%%%===================================================================
%% Set connection manager PID (sent by facade after init)
handle_cast({set_connection_manager_pid, Pid}, State) ->
{noreply, State#state{connection_manager_pid = Pid}};
handle_cast({do_publish, PublishMsg, Qos, BinaryTopic, Payload, Opts, MsgId}, State) ->
%% HYBRID: Send to gateway for routing AND discover via DHT
%% This enables both gateway-centric and pure P2P topologies
?LOG_DEBUG("[~s] Publishing message to topic ~s (qos=~p, msg_id=~s) via gateway + DHT",
[State#state.node_id, BinaryTopic, Qos, MsgId]),
%% Handle QoS 1 (at-least-once delivery) - delegate to QoS module
{ok, UpdatedPendingPubacks} = macula_pubsub_qos:track_message(
MsgId, BinaryTopic, Payload, Qos, State#state.pending_pubacks
),
State2 = State#state{pending_pubacks = UpdatedPendingPubacks},
%% Send publish to gateway for routing to connected subscribers
%% This is essential for gateway-centric topologies where peers can't reach each other
case State#state.connection_manager_pid of
undefined ->
?LOG_WARNING("[~s] No connection manager - cannot send publish to gateway",
[State#state.node_id]);
ConnMgrPid ->
%% Send publish message via the gateway connection
macula_connection:send_message(ConnMgrPid, publish, PublishMsg),
?LOG_DEBUG("[~s] Sent publish to gateway for topic ~s",
[State#state.node_id, BinaryTopic])
end,
%% Also discover subscribers via DHT for pure P2P routing (if available)
gen_server:cast(self(), {discover_subscribers, BinaryTopic, Payload, Qos, Opts}),
{noreply, State2};
%% @private
%% Handle async discovery of remote subscribers (mesh-wide pub/sub)
%% Delegate to DHT module for discovery and routing
handle_cast({discover_subscribers, Topic, Payload, Qos, _Opts}, State) ->
%% Delegate to DHT module for discovery (async)
case macula_pubsub_dht:discover_subscribers(
Topic, Payload, Qos,
State#state.connection_manager_pid,
State#state.service_registry,
State#state.msg_id_counter
) of
{cached, Subscribers, UpdatedRegistry} ->
%% Cache hit - route to subscribers immediately
macula_pubsub_dht:route_to_subscribers(
Topic, Payload, Qos, Subscribers, State#state.node_id
),
{noreply, State#state{service_registry = UpdatedRegistry}};
{query_sent, PendingQueries, MsgId, UpdatedRegistry} ->
%% Query sent - track it
UpdatedPendingQueries = maps:merge(
State#state.pending_subscriber_queries,
PendingQueries
),
{_MsgId, NewCounter} = macula_utils:next_message_id(State#state.msg_id_counter),
State2 = State#state{
pending_subscriber_queries = UpdatedPendingQueries,
service_registry = UpdatedRegistry,
msg_id_counter = NewCounter
},
?LOG_DEBUG("[~s] DHT query sent for topic ~s (MsgId: ~s)",
[State#state.node_id, Topic, MsgId]),
{noreply, State2}
end;
%%%===================================================================
%%% Incoming Publish Routing
%%%===================================================================
handle_cast({incoming_publish, Msg}, State) ->
%% Handle incoming publish (for subscriptions)
%% Support both atom and binary keys from MessagePack decoding
Topic = extract_topic(Msg),
Payload = extract_payload(Msg),
?LOG_DEBUG("[~s] Received PUBLISH message: topic=~s, payload_size=~p bytes",
[State#state.node_id, Topic, byte_size(Payload)]),
%% Find matching subscriptions (delegate to subscription module)
Config = #{
topic_separator => State#state.topic_separator,
topic_wildcard_single => State#state.topic_wildcard_single,
topic_wildcard_multi => State#state.topic_wildcard_multi
},
SubscriptionMatches = macula_pubsub_subscription:find_matches(
Topic, State#state.subscriptions, Config
),
%% Invoke callbacks for matching subscriptions (delegate to subscription module)
macula_pubsub_subscription:invoke_callbacks(
SubscriptionMatches, Topic, Payload, State#state.node_id
),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%%%===================================================================
%%% Timer Handlers
%%%===================================================================
handle_info({resubscribe, Topic}, State) ->
case maps:get(Topic, State#state.advertised_subscriptions, undefined) of
undefined ->
%% Subscription was removed, don't re-advertise
?LOG_DEBUG("[~s] Skipping re-subscription for ~s (no longer subscribed)",
[State#state.node_id, Topic]),
{noreply, State};
SubInfo ->
%% Cancel old timer
OldTimerRef = maps:get(timer_ref, SubInfo),
erlang:cancel_timer(OldTimerRef),
%% Re-advertise (delegate to DHT module)
SubRef = maps:get(sub_ref, SubInfo),
{ok, UpdatedSubInfo} = macula_pubsub_dht:advertise_subscription(
Topic, SubRef, State#state.node_id, State#state.url,
State#state.connection_manager_pid
),
%% Update advertised subscriptions map
UpdatedAdvertised = (State#state.advertised_subscriptions)#{Topic => UpdatedSubInfo},
State2 = State#state{advertised_subscriptions = UpdatedAdvertised},
?LOG_DEBUG("[~s] Re-advertised subscription for topic ~s",
[State#state.node_id, Topic]),
{noreply, State2}
end;
handle_info({puback_timeout, MsgId}, State) ->
%% Delegate QoS timeout handling to QoS module
case macula_pubsub_qos:handle_timeout(MsgId, State#state.connection_manager_pid, State#state.pending_pubacks) of
{retry, UpdatedPending, PublishMsg} ->
%% Send retry message
case macula_connection:send_message(State#state.connection_manager_pid, publish, PublishMsg) of
ok ->
{noreply, State#state{pending_pubacks = UpdatedPending}};
{error, SendError} ->
?LOG_ERROR("[~s] Failed to retry publish: ~p - giving up", [State#state.node_id, SendError]),
%% Remove from pending on send failure
{noreply, State#state{pending_pubacks = maps:remove(MsgId, UpdatedPending)}}
end;
{give_up, UpdatedPending} ->
{noreply, State#state{pending_pubacks = UpdatedPending}};
{not_found, UpdatedPending} ->
{noreply, State#state{pending_pubacks = UpdatedPending}}
end;
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
?LOG_INFO("PubSub handler terminating"),
ok.
%%%===================================================================
%%% Internal functions - Helpers
%%%===================================================================
%% @doc Generate next message ID
-spec next_message_id(#state{}) -> {binary(), #state{}}.
next_message_id(State) ->
Counter = State#state.msg_id_counter,
{MsgId, NewCounter} = macula_utils:next_message_id(Counter),
{MsgId, State#state{msg_id_counter = NewCounter}}.
%% @doc Ensure value is binary
-spec ensure_binary(binary() | list() | atom()) -> binary().
ensure_binary(Value) ->
macula_utils:ensure_binary(Value).
%% @doc Encode map/list to JSON binary
-spec encode_json(map() | list()) -> binary().
encode_json(Data) ->
macula_utils:encode_json(Data).
%% @doc Generate a random node ID
-spec generate_node_id() -> binary().
generate_node_id() ->
macula_utils:generate_node_id().
%% @doc Encode payload data to binary (pattern matching on type)
-spec encode_payload(binary() | map() | list()) -> binary().
encode_payload(Data) when is_binary(Data) -> Data;
encode_payload(Data) when is_map(Data) -> encode_json(Data);
encode_payload(Data) when is_list(Data) -> list_to_binary(Data).
%% @doc Extract topic from message, supporting both atom and binary keys.
%% MessagePack decoding may use either key format.
-spec extract_topic(map()) -> binary().
extract_topic(Msg) ->
case maps:get(topic, Msg, undefined) of
undefined -> maps:get(<<"topic">>, Msg);
Topic -> Topic
end.
%% @doc Extract payload from message, supporting both atom and binary keys.
%% MessagePack decoding may use either key format.
-spec extract_payload(map()) -> binary().
extract_payload(Msg) ->
case maps:get(payload, Msg, undefined) of
undefined -> maps:get(<<"payload">>, Msg);
Payload -> Payload
end.