Packages
macula
0.48.2
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_subscription.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Subscription management for pub/sub.
%%%
%%% Responsibilities:
%%% - Store and retrieve subscriptions
%%% - Pattern matching with wildcards (*, **)
%%% - Find matching subscriptions for a topic
%%% - Invoke subscriber callbacks
%%%
%%% Extracted from macula_pubsub_handler.erl (Phase 4)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_subscription).
-include_lib("kernel/include/logger.hrl").
%% API
-export([add_subscription/4, remove_subscription/2,
find_matches/3, invoke_callbacks/4, invoke_callbacks/5]).
-type topic() :: binary().
-type callback() :: fun((map()) -> ok).
-type subscription_ref() :: reference().
-type node_id() :: binary().
-type payload() :: binary().
-type subscriptions() :: #{subscription_ref() => {topic(), callback()}}.
-export_type([subscriptions/0]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Add a subscription.
%% Returns {ok, UpdatedSubscriptions, SubRef}.
-spec add_subscription(topic(), callback(), subscriptions(), subscription_ref()) ->
{ok, subscriptions(), subscription_ref()}.
add_subscription(Topic, Callback, Subscriptions, SubRef) ->
UpdatedSubscriptions = Subscriptions#{SubRef => {Topic, Callback}},
{ok, UpdatedSubscriptions, SubRef}.
%% @doc Remove a subscription.
%% Returns {ok, UpdatedSubscriptions, Topic} | {error, not_found}.
-spec remove_subscription(subscription_ref(), subscriptions()) ->
{ok, subscriptions(), topic()} | {error, not_found}.
remove_subscription(SubRef, Subscriptions) ->
SubInfo = maps:get(SubRef, Subscriptions, undefined),
do_remove_subscription(SubInfo, SubRef, Subscriptions).
%% @private Subscription not found
do_remove_subscription(undefined, _SubRef, _Subscriptions) ->
{error, not_found};
%% @private Remove subscription
do_remove_subscription({Topic, _Callback}, SubRef, Subscriptions) ->
UpdatedSubscriptions = maps:remove(SubRef, Subscriptions),
{ok, UpdatedSubscriptions, Topic}.
%% @doc Find matching subscriptions for a topic.
%% Returns list of {SubRef, {Pattern, Callback}} tuples.
-spec find_matches(topic(), subscriptions(), #{atom() => binary()}) ->
[{subscription_ref(), {topic(), callback()}}].
find_matches(Topic, Subscriptions, Config) ->
Separator = maps:get(topic_separator, Config),
WildcardSingle = maps:get(topic_wildcard_single, Config),
WildcardMulti = maps:get(topic_wildcard_multi, Config),
lists:filter(
fun({_SubRef, {Pattern, _Callback}}) ->
macula_utils:topic_matches(Pattern, Topic, Separator, WildcardSingle, WildcardMulti)
end,
maps:to_list(Subscriptions)
).
%% @doc Invoke callbacks for matching subscriptions.
%% Spawns async tasks to invoke each callback.
-spec invoke_callbacks([{subscription_ref(), {topic(), callback()}}], topic(), payload(), node_id()) -> ok.
invoke_callbacks(Matches, Topic, Payload, NodeId) ->
invoke_callbacks(Matches, Topic, Payload, NodeId, undefined).
-spec invoke_callbacks([{subscription_ref(), {topic(), callback()}}], topic(), payload(), node_id(), list() | undefined) -> ok.
invoke_callbacks([], _Topic, _Payload, _NodeId, _Trace) ->
ok;
invoke_callbacks(Matches, Topic, Payload, NodeId, Trace) ->
?LOG_INFO("[~s] Found ~p subscription(s) for topic: ~s",
[NodeId, length(Matches), Topic]),
%% Invoke all matching callbacks asynchronously
lists:foreach(
fun({_SubRef, {SubTopic, Callback}}) ->
spawn(fun() ->
invoke_single_callback(Callback, Topic, SubTopic, Payload, NodeId, Trace)
end)
end,
Matches
),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private Invoke a single callback with payload decoding
invoke_single_callback(Callback, Topic, SubTopic, Payload, NodeId, Trace) ->
DecodedPayload = safe_decode_json(Payload),
Base = #{
topic => Topic,
matched_pattern => SubTopic,
payload => DecodedPayload
},
PublishData = case Trace of
undefined -> Base;
_ -> Base#{'_trace' => Trace}
end,
handle_callback_invocation(catch Callback(PublishData), Topic, NodeId).
%% @private Safely decode JSON payload
safe_decode_json(Payload) ->
handle_decode_result(catch macula_utils:decode_json(Payload), Payload).
%% @private Handle JSON decode result
handle_decode_result({'EXIT', _}, Original) ->
Original;
handle_decode_result(Decoded, _Original) ->
Decoded.
%% @private Handle callback invocation result
handle_callback_invocation({'EXIT', {Reason, Stacktrace}}, Topic, NodeId) ->
?LOG_ERROR("[~s] Callback error for topic ~s: ~p~nStack: ~p",
[NodeId, Topic, Reason, Stacktrace]);
handle_callback_invocation({'EXIT', Reason}, Topic, NodeId) ->
?LOG_ERROR("[~s] Callback error for topic ~s: ~p", [NodeId, Topic, Reason]);
handle_callback_invocation(_Result, Topic, NodeId) ->
?LOG_DEBUG("[~s] Invoked callback for topic ~s", [NodeId, Topic]).