Packages
macula
0.8.5
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]).
-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) ->
case maps:get(SubRef, Subscriptions, undefined) of
undefined ->
{error, not_found};
{Topic, _Callback} ->
UpdatedSubscriptions = maps:remove(SubRef, Subscriptions),
{ok, UpdatedSubscriptions, Topic}
end.
%% @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([], _Topic, _Payload, _NodeId) ->
ok;
invoke_callbacks(Matches, Topic, Payload, NodeId) ->
?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() ->
try
%% Decode payload if it's JSON
DecodedPayload = try
macula_utils:decode_json(Payload)
catch
_:_ -> Payload
end,
PublishData = #{
topic => Topic,
matched_pattern => SubTopic,
payload => DecodedPayload
},
Callback(PublishData),
?LOG_DEBUG("[~s] Invoked callback for topic ~s", [NodeId, Topic])
catch
Error:Reason:Stack ->
?LOG_ERROR("[~s] Callback error for topic ~s: ~p:~p~nStack: ~p",
[NodeId, Topic, Error, Reason, Stack])
end
end)
end,
Matches
),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================