Packages
macula
0.6.7
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_registry.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Local subscription registry for pub/sub.
%%% Maps topic patterns to local subscribers (callback PIDs).
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_registry).
%% API
-export([
new/0,
subscribe/4,
unsubscribe/3,
match/2,
list_patterns/1,
get_subscription/3,
size/1
]).
%% Types
-type subscription() :: #{
subscriber_id := binary(),
pattern := binary(),
callback := pid()
}.
-type registry() :: #{
subscriptions := [subscription()],
pattern_index := #{binary() => [subscription()]} % Pattern -> [Subscriptions]
}.
-export_type([subscription/0, registry/0]).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Create new empty registry.
-spec new() -> registry().
new() ->
#{
subscriptions => [],
pattern_index => #{}
}.
%% @doc Subscribe to a pattern.
%% If subscription already exists (same subscriber_id + pattern), updates callback.
-spec subscribe(registry(), binary(), binary(), pid()) -> registry().
subscribe(#{subscriptions := Subs, pattern_index := Index} = Registry, SubscriberId, Pattern, Callback) ->
%% Create subscription
Subscription = #{
subscriber_id => SubscriberId,
pattern => Pattern,
callback => Callback
},
%% Check if subscription already exists
case find_subscription(Subs, SubscriberId, Pattern) of
{found, _OldSub} ->
%% Update existing subscription
NewSubs = update_subscription(Subs, Subscription),
NewIndex = update_pattern_index(Index, Pattern, NewSubs),
Registry#{subscriptions => NewSubs, pattern_index => NewIndex};
not_found ->
%% Add new subscription
NewSubs = [Subscription | Subs],
NewIndex = add_to_pattern_index(Index, Pattern, Subscription),
Registry#{subscriptions => NewSubs, pattern_index => NewIndex}
end.
%% @doc Unsubscribe from a pattern.
-spec unsubscribe(registry(), binary(), binary()) -> registry().
unsubscribe(#{subscriptions := Subs, pattern_index := Index} = Registry, SubscriberId, Pattern) ->
case find_subscription(Subs, SubscriberId, Pattern) of
{found, Sub} ->
%% Remove subscription
NewSubs = lists:filter(
fun(S) ->
not (maps:get(subscriber_id, S) =:= SubscriberId andalso
maps:get(pattern, S) =:= Pattern)
end,
Subs
),
NewIndex = remove_from_pattern_index(Index, Pattern, Sub),
Registry#{subscriptions => NewSubs, pattern_index => NewIndex};
not_found ->
Registry % No change
end.
%% @doc Find subscriptions matching a topic.
-spec match(registry(), binary()) -> [subscription()].
match(#{subscriptions := Subs}, Topic) ->
lists:filter(
fun(Sub) ->
Pattern = maps:get(pattern, Sub),
macula_pubsub_topic:matches(Topic, Pattern)
end,
Subs
).
%% @doc List all unique patterns.
-spec list_patterns(registry()) -> [binary()].
list_patterns(#{pattern_index := Index}) ->
maps:keys(Index).
%% @doc Get specific subscription.
-spec get_subscription(registry(), binary(), binary()) -> {ok, subscription()} | not_found.
get_subscription(#{subscriptions := Subs}, SubscriberId, Pattern) ->
case find_subscription(Subs, SubscriberId, Pattern) of
{found, Sub} -> {ok, Sub};
not_found -> not_found
end.
%% @doc Get number of subscriptions.
-spec size(registry()) -> non_neg_integer().
size(#{subscriptions := Subs}) ->
length(Subs).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Find subscription by subscriber_id and pattern.
-spec find_subscription([subscription()], binary(), binary()) -> {found, subscription()} | not_found.
find_subscription([], _SubscriberId, _Pattern) ->
not_found;
find_subscription([Sub | Rest], SubscriberId, Pattern) ->
case maps:get(subscriber_id, Sub) =:= SubscriberId andalso
maps:get(pattern, Sub) =:= Pattern of
true -> {found, Sub};
false -> find_subscription(Rest, SubscriberId, Pattern)
end.
%% @doc Update existing subscription in list.
-spec update_subscription([subscription()], subscription()) -> [subscription()].
update_subscription(Subs, NewSub) ->
SubscriberId = maps:get(subscriber_id, NewSub),
Pattern = maps:get(pattern, NewSub),
lists:map(
fun(Sub) ->
case maps:get(subscriber_id, Sub) =:= SubscriberId andalso
maps:get(pattern, Sub) =:= Pattern of
true -> NewSub;
false -> Sub
end
end,
Subs
).
%% @doc Add subscription to pattern index.
-spec add_to_pattern_index(#{binary() => [subscription()]}, binary(), subscription()) ->
#{binary() => [subscription()]}.
add_to_pattern_index(Index, Pattern, Subscription) ->
Existing = maps:get(Pattern, Index, []),
Index#{Pattern => [Subscription | Existing]}.
%% @doc Update pattern index after subscription update.
-spec update_pattern_index(#{binary() => [subscription()]}, binary(), [subscription()]) ->
#{binary() => [subscription()]}.
update_pattern_index(Index, Pattern, AllSubs) ->
%% Rebuild index entry for this pattern
PatternSubs = lists:filter(
fun(Sub) -> maps:get(pattern, Sub) =:= Pattern end,
AllSubs
),
Index#{Pattern => PatternSubs}.
%% @doc Remove subscription from pattern index.
-spec remove_from_pattern_index(#{binary() => [subscription()]}, binary(), subscription()) ->
#{binary() => [subscription()]}.
remove_from_pattern_index(Index, Pattern, Subscription) ->
Existing = maps:get(Pattern, Index, []),
SubscriberId = maps:get(subscriber_id, Subscription),
NewList = lists:filter(
fun(Sub) -> maps:get(subscriber_id, Sub) =/= SubscriberId end,
Existing
),
case NewList of
[] ->
%% No more subscriptions for this pattern, remove key
maps:remove(Pattern, Index);
_ ->
Index#{Pattern => NewList}
end.