Packages

macula

0.20.21
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_qos.erl
Raw

src/macula_pubsub_system/macula_pubsub_qos.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% QoS (Quality of Service) manager for pub/sub.
%%%
%%% Handles QoS 1 (at-least-once delivery) logic:
%%% - Message tracking with timeout timers
%%% - Automatic retry on timeout (up to max retries)
%%% - Acknowledgment handling
%%%
%%% Extracted from macula_pubsub_handler.erl (Phase 2)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_qos).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
%% API
-export([track_message/5, handle_timeout/3, handle_ack/2, get_pending/1]).
-type message_id() :: binary().
-type topic() :: binary().
-type payload() :: binary().
-type qos() :: 0 | 1.
-type retry_count() :: non_neg_integer().
-type timer_ref() :: reference().
-type connection_manager_pid() :: pid().
-type pending_pubacks() :: #{message_id() => {topic(), payload(), qos(), retry_count(), timer_ref()}}.
-export_type([pending_pubacks/0]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Track a message for QoS 1 acknowledgment.
%% Starts a timeout timer and stores message in pending map.
%% Returns updated pending_pubacks map.
-spec track_message(message_id(), topic(), payload(), qos(), pending_pubacks()) ->
{ok, pending_pubacks()} | {error, term()}.
track_message(_MsgId, _Topic, _Payload, Qos, PendingPubacks) when Qos =/= 1 ->
%% Only track QoS 1 messages
{ok, PendingPubacks};
track_message(MsgId, Topic, Payload, Qos, PendingPubacks) when Qos =:= 1 ->
%% Start timeout timer
TimerRef = erlang:send_after(?PUBACK_TIMEOUT, self(), {puback_timeout, MsgId}),
%% Store pending acknowledgment with retry count 0
UpdatedPending = PendingPubacks#{MsgId => {Topic, Payload, Qos, 0, TimerRef}},
{ok, UpdatedPending}.
%% @doc Handle timeout for a pending message.
%% Retries sending if under max retries, otherwise gives up.
%% Returns {retry, UpdatedPending, PublishMsg} | {give_up, UpdatedPending}.
-spec handle_timeout(message_id(), connection_manager_pid(), pending_pubacks()) ->
{retry, pending_pubacks(), map()} | {give_up, pending_pubacks()} | {not_found, pending_pubacks()}.
handle_timeout(MsgId, _ConnMgrPid, PendingPubacks) ->
PendingInfo = maps:get(MsgId, PendingPubacks, undefined),
do_handle_timeout(PendingInfo, MsgId, PendingPubacks).
%% @private Message not found (already acknowledged or removed)
do_handle_timeout(undefined, _MsgId, PendingPubacks) ->
{not_found, PendingPubacks};
%% @private Message found - check retry count
do_handle_timeout({Topic, Payload, Qos, RetryCount, _OldTimerRef}, MsgId, PendingPubacks) ->
NewRetryCount = RetryCount + 1,
MaxRetriesReached = NewRetryCount >= ?PUBACK_MAX_RETRIES,
handle_timeout_retry(MaxRetriesReached, MsgId, Topic, Payload, Qos, NewRetryCount, PendingPubacks).
%% @private Max retries reached, give up
handle_timeout_retry(true, MsgId, Topic, _Payload, _Qos, NewRetryCount, PendingPubacks) ->
?LOG_ERROR("[QoS] PUBACK timeout for message ~p (topic: ~s) after ~p retries - giving up",
[MsgId, Topic, NewRetryCount]),
{give_up, maps:remove(MsgId, PendingPubacks)};
%% @private Retry sending the message
handle_timeout_retry(false, MsgId, Topic, Payload, Qos, NewRetryCount, PendingPubacks) ->
?LOG_WARNING("[QoS] PUBACK timeout for message ~p (topic: ~s) - retry ~p/~p",
[MsgId, Topic, NewRetryCount, ?PUBACK_MAX_RETRIES]),
PublishMsg = #{
topic => Topic,
payload => Payload,
qos => Qos,
retain => false,
message_id => MsgId
},
NewTimerRef = erlang:send_after(?PUBACK_TIMEOUT, self(), {puback_timeout, MsgId}),
UpdatedPending = PendingPubacks#{
MsgId => {Topic, Payload, Qos, NewRetryCount, NewTimerRef}
},
{retry, UpdatedPending, PublishMsg}.
%% @doc Handle acknowledgment for a message.
%% Cancels timer and removes message from pending map.
%% Returns updated pending_pubacks map.
-spec handle_ack(message_id(), pending_pubacks()) -> pending_pubacks().
handle_ack(MsgId, PendingPubacks) ->
PendingInfo = maps:get(MsgId, PendingPubacks, undefined),
do_handle_ack(PendingInfo, MsgId, PendingPubacks).
%% @private Message not found, return unchanged
do_handle_ack(undefined, _MsgId, PendingPubacks) ->
PendingPubacks;
%% @private Cancel timer and remove message
do_handle_ack({_Topic, _Payload, _Qos, _RetryCount, TimerRef}, MsgId, PendingPubacks) ->
erlang:cancel_timer(TimerRef),
maps:remove(MsgId, PendingPubacks).
%% @doc Get list of pending message IDs (for testing/debugging).
-spec get_pending(pending_pubacks()) -> [message_id()].
get_pending(PendingPubacks) ->
maps:keys(PendingPubacks).
%%%===================================================================
%%% Internal functions
%%%===================================================================