Packages

macula

0.20.2
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) ->
%% Use cast to avoid blocking on publish operations
%% The internal handler will manage async delivery and QoS retries
gen_server:cast(Pid, {publish_async, Topic, Data, Opts}),
ok.
-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">>),
PeerId = maps:get(peer_id, Opts, erlang:unique_integer([monotonic, positive])),
%% Register in gproc for incoming publish message routing
%% Use {Realm, PeerId} to support multiple peer connections per realm
true = gproc:reg({n, l, {pubsub_handler, Realm, PeerId}}),
?LOG_INFO("PubSub handler registered in gproc for realm ~s, peer_id ~p", [Realm, PeerId]),
%% 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),
%% Store subscription locally first (delegate to subscription module)
{ok, UpdatedSubscriptions, SubRef} = macula_pubsub_subscription:add_subscription(
BinaryTopic, Callback, State#state.subscriptions, SubRef
),
State2 = State#state{subscriptions = UpdatedSubscriptions},
%% Send subscribe message via connection manager ASYNC (fire-and-forget)
%% This prevents blocking when the connection is busy or slow
SubscribeMsg = #{
topics => [BinaryTopic],
qos => 0
},
macula_connection:send_message_async(State#state.connection_manager_pid, subscribe, SubscribeMsg),
%% 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};
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 ASYNC (fire-and-forget)
%% This prevents blocking when the connection is busy or slow
UnsubscribeMsg = #{
topics => [Topic]
},
macula_connection:send_message_async(State#state.connection_manager_pid, unsubscribe, UnsubscribeMsg),
%% 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}
end;
%%%===================================================================
%%% Publish
%%%===================================================================
handle_call({publish, Topic, Data, Opts}, _From, State) ->
%% Check if we have a connection manager and are connected
do_sync_publish(State#state.connection_manager_pid, Topic, Data, Opts, State);
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}};
%% Async publish (fire-and-forget from caller's perspective)
%% This handles the {publish_async, ...} cast from the publish/4 API function
handle_cast({publish_async, Topic, Data, Opts}, State) ->
?LOG_INFO("[PubSubHandler] publish_async received: topic=~p", [Topic]),
do_async_publish(State#state.connection_manager_pid, Topic, Data, Opts, State);
handle_cast({do_publish, PublishMsg, Qos, BinaryTopic, Payload, _Opts, MsgId}, State) ->
%% Send publish to gateway for routing - gateway handles DHT lookup and subscriber delivery
%% This avoids duplicate delivery (previously HYBRID mode sent via BOTH paths)
?LOG_INFO("[PubSubHandler] do_publish: topic=~s, ConnMgr=~p",
[BinaryTopic, State#state.connection_manager_pid]),
%% 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
%% Gateway will lookup subscribers via DHT and route messages accordingly
send_publish_to_gateway(State#state.connection_manager_pid, PublishMsg, BinaryTopic),
{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) ->
SubInfo = maps:get(Topic, State#state.advertised_subscriptions, undefined),
do_resubscribe(SubInfo, Topic, State);
handle_info({puback_timeout, MsgId}, State) ->
%% Delegate QoS timeout handling to QoS module
TimeoutResult = macula_pubsub_qos:handle_timeout(MsgId, State#state.connection_manager_pid, State#state.pending_pubacks),
handle_puback_timeout_result(TimeoutResult, State);
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(#{topic := Topic}) -> Topic;
extract_topic(#{<<"topic">> := Topic}) -> Topic.
%% @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(#{payload := Payload}) -> Payload;
extract_payload(#{<<"payload">> := Payload}) -> Payload.
%%%===================================================================
%%% Publish helpers
%%%===================================================================
%% @private No connection manager - cannot publish
do_sync_publish(undefined, _Topic, _Data, _Opts, State) ->
{reply, {error, not_connected}, State};
%% @private Connection manager available - check status
do_sync_publish(ConnMgrPid, Topic, Data, Opts, State) ->
Status = macula_connection:get_status(ConnMgrPid),
do_sync_publish_with_status(Status, Topic, Data, Opts, State).
%% @private Not connected
do_sync_publish_with_status(Status, _Topic, _Data, _Opts, State) when Status =/= connected ->
{reply, {error, not_connected}, State};
%% @private Connected - build and send publish message
do_sync_publish_with_status(connected, Topic, Data, Opts, State) ->
Qos = maps:get(qos, Opts, 0),
Retain = maps:get(retain, Opts, false),
{MsgId, State2} = next_message_id(State),
BinaryTopic = ensure_binary(Topic),
Payload = encode_payload(Data),
PublishMsg = #{
topic => BinaryTopic,
payload => Payload,
qos => Qos,
retain => Retain,
message_id => MsgId
},
%% Send publish message asynchronously using cast to self
gen_server:cast(self(), {do_publish, PublishMsg, Qos, BinaryTopic, Payload, Opts, MsgId}),
{reply, ok, State2}.
%% @private No connection manager - silently drop (fire-and-forget semantics)
do_async_publish(undefined, _Topic, _Data, _Opts, State) ->
?LOG_WARNING("[PubSubHandler] Publish dropped - no connection manager"),
{noreply, State};
%% @private Connection manager available - check authorization then send async
do_async_publish(_ConnMgrPid, Topic, Data, Opts, State) ->
BinaryTopic = ensure_binary(Topic),
%% Authorization check (v0.17.0+)
CallerDID = get_caller_did(State, Opts),
UcanToken = maps:get(ucan_token, Opts, undefined),
case macula_authorization:check_publish(CallerDID, BinaryTopic, UcanToken, Opts) of
{ok, authorized} ->
Qos = maps:get(qos, Opts, 0),
Retain = maps:get(retain, Opts, false),
{MsgId, State2} = next_message_id(State),
Payload = encode_payload(Data),
PublishMsg = #{
topic => BinaryTopic,
payload => Payload,
qos => Qos,
retain => Retain,
message_id => MsgId
},
?LOG_INFO("[PubSubHandler] Sending to do_publish: topic=~s", [BinaryTopic]),
gen_server:cast(self(), {do_publish, PublishMsg, Qos, BinaryTopic, Payload, Opts, MsgId}),
{noreply, State2};
{error, Reason} ->
?LOG_WARNING("[~s] Publish to ~s denied: ~p (caller: ~s)",
[State#state.node_id, BinaryTopic, Reason, CallerDID]),
{noreply, State}
end.
%% @private No connection manager - log error
send_publish_to_gateway(undefined, _PublishMsg, _BinaryTopic) ->
?LOG_ERROR("[PubSubHandler] No connection manager for publish!");
%% @private Send publish message via the gateway connection ASYNC (fire-and-forget)
send_publish_to_gateway(ConnMgrPid, PublishMsg, BinaryTopic) ->
?LOG_INFO("[PubSubHandler] Calling send_message_async: pid=~p, topic=~s",
[ConnMgrPid, BinaryTopic]),
macula_connection:send_message_async(ConnMgrPid, publish, PublishMsg),
?LOG_INFO("[PubSubHandler] send_message_async returned").
%%%===================================================================
%%% Resubscription helpers
%%%===================================================================
%% @private Subscription was removed, don't re-advertise
do_resubscribe(undefined, Topic, State) ->
?LOG_DEBUG("[~s] Skipping re-subscription for ~s (no longer subscribed)",
[State#state.node_id, Topic]),
{noreply, State};
%% @private Re-advertise subscription
do_resubscribe(SubInfo, Topic, State) ->
%% 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}.
%%%===================================================================
%%% QoS timeout helpers
%%%===================================================================
%% @private Handle retry case - send retry message ASYNC
handle_puback_timeout_result({retry, UpdatedPending, PublishMsg}, State) ->
macula_connection:send_message_async(State#state.connection_manager_pid, publish, PublishMsg),
{noreply, State#state{pending_pubacks = UpdatedPending}};
%% @private Handle give up case
handle_puback_timeout_result({give_up, UpdatedPending}, State) ->
{noreply, State#state{pending_pubacks = UpdatedPending}};
%% @private Handle not found case
handle_puback_timeout_result({not_found, UpdatedPending}, State) ->
{noreply, State#state{pending_pubacks = UpdatedPending}}.
%%%===================================================================
%%% Authorization Helpers (v0.17.0+)
%%%===================================================================
%% @doc Get caller DID from options or derive from local identity.
%% The caller_did can be provided explicitly in Opts (from TLS cert extraction),
%% or we fall back to the local node's realm identity.
-spec get_caller_did(#state{}, map()) -> binary().
get_caller_did(State, Opts) ->
case maps:get(caller_did, Opts, undefined) of
undefined ->
%% Fallback: derive DID from realm (e.g., "did:macula:io.macula.default")
Realm = State#state.realm,
<<"did:macula:", Realm/binary>>;
CallerDID when is_binary(CallerDID) ->
CallerDID
end.