Packages

macula

0.7.28
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_gateway_pubsub.erl
Raw

src/macula_gateway_pubsub.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Pub/Sub Handler GenServer - manages topic subscriptions and message routing.
%%%
%%% Responsibilities:
%%% - Subscribe/unsubscribe streams to topics
%%% - Route published messages to matching subscribers
%%% - Support wildcard topics (* single-level, ** multi-level)
%%% - Track bidirectional mapping (topic ↔ stream)
%%% - Monitor stream processes for automatic cleanup
%%%
%%% Extracted from macula_gateway.erl (Phase 3)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_pubsub).
-behaviour(gen_server).
%% API
-export([
start_link/1,
stop/1,
subscribe/3,
unsubscribe/3,
publish/3,
get_subscribers/2,
get_stream_topics/2
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
opts :: map(),
subscriptions :: #{binary() => [pid()]}, % topic => [stream_pids]
stream_subscriptions :: #{pid() => [binary()]}, % stream_pid => [topics]
monitors :: #{reference() => {pid(), binary()}} % monitor_ref => {stream_pid, topic}
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the pub/sub handler with options.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Stop the pub/sub handler.
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Subscribe a stream to a topic (supports wildcards).
-spec subscribe(pid(), pid(), binary()) -> ok.
subscribe(Pid, Stream, Topic) ->
gen_server:call(Pid, {subscribe, Stream, Topic}).
%% @doc Unsubscribe a stream from a topic.
-spec unsubscribe(pid(), pid(), binary()) -> ok.
unsubscribe(Pid, Stream, Topic) ->
gen_server:call(Pid, {unsubscribe, Stream, Topic}).
%% @doc Publish a message to a topic (routes to matching subscribers).
-spec publish(pid(), binary(), map()) -> ok.
publish(Pid, Topic, Payload) ->
gen_server:call(Pid, {publish, Topic, Payload}).
%% @doc Get all subscribers for a topic (exact and wildcard matches).
-spec get_subscribers(pid(), binary()) -> {ok, [pid()]}.
get_subscribers(Pid, Topic) ->
gen_server:call(Pid, {get_subscribers, Topic}).
%% @doc Get all topics a stream is subscribed to.
-spec get_stream_topics(pid(), pid()) -> {ok, [binary()]}.
get_stream_topics(Pid, Stream) ->
gen_server:call(Pid, {get_stream_topics, Stream}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
io:format("[PubSub] Initializing pub/sub handler~n"),
State = #state{
opts = Opts,
subscriptions = #{},
stream_subscriptions = #{},
monitors = #{}
},
io:format("[PubSub] Pub/sub handler initialized~n"),
{ok, State}.
handle_call({subscribe, Stream, Topic}, _From, State) when is_pid(Stream), is_binary(Topic) ->
io:format("[PubSub ~p] SUBSCRIBE called: Stream=~p, Topic=~s~n", [self(), Stream, Topic]),
%% Check if already subscribed
CurrentTopics = maps:get(Stream, State#state.stream_subscriptions, []),
NewState = case lists:member(Topic, CurrentTopics) of
true ->
io:format("[PubSub ~p] Stream ~p already subscribed to ~s~n", [self(), Stream, Topic]),
%% Already subscribed - idempotent
State;
false ->
%% Add to topic → streams mapping
Subscribers = maps:get(Topic, State#state.subscriptions, []),
NewSubscriptions = maps:put(Topic, [Stream | Subscribers], State#state.subscriptions),
io:format("[PubSub ~p] Added stream ~p to topic ~s (total subscribers: ~p)~n",
[self(), Stream, Topic, length([Stream | Subscribers])]),
%% Add to stream → topics mapping
NewTopics = [Topic | CurrentTopics],
NewStreamSubs = maps:put(Stream, NewTopics, State#state.stream_subscriptions),
%% Monitor stream if first subscription
NewMonitors = case length(CurrentTopics) of
0 ->
MonitorRef = erlang:monitor(process, Stream),
io:format("[PubSub ~p] Monitoring stream ~p (first subscription)~n", [self(), Stream]),
maps:put(MonitorRef, {Stream, Topic}, State#state.monitors);
_ ->
State#state.monitors
end,
io:format("[PubSub ~p] Current subscriptions map: ~p~n", [self(), NewSubscriptions]),
State#state{
subscriptions = NewSubscriptions,
stream_subscriptions = NewStreamSubs,
monitors = NewMonitors
}
end,
{reply, ok, NewState};
handle_call({unsubscribe, Stream, Topic}, _From, State) when is_pid(Stream), is_binary(Topic) ->
%% Remove from stream → topics mapping
CurrentTopics = maps:get(Stream, State#state.stream_subscriptions, []),
NewTopics = lists:delete(Topic, CurrentTopics),
NewStreamSubs = case NewTopics of
[] -> maps:remove(Stream, State#state.stream_subscriptions);
_ -> maps:put(Stream, NewTopics, State#state.stream_subscriptions)
end,
%% Remove from topic → streams mapping
Subscribers = maps:get(Topic, State#state.subscriptions, []),
NewSubscribers = lists:delete(Stream, Subscribers),
NewSubscriptions = case NewSubscribers of
[] -> maps:remove(Topic, State#state.subscriptions);
_ -> maps:put(Topic, NewSubscribers, State#state.subscriptions)
end,
NewState = State#state{
subscriptions = NewSubscriptions,
stream_subscriptions = NewStreamSubs
},
{reply, ok, NewState};
handle_call({publish, Topic, Payload}, _From, State) when is_binary(Topic) ->
%% Find all matching subscribers (exact + wildcard patterns)
MatchingStreams = find_matching_subscribers(Topic, State),
%% Send message to all matching streams
lists:foreach(fun(Stream) ->
case erlang:is_process_alive(Stream) of
true ->
Stream ! {publish, Topic, Payload};
false ->
ok %% Ignore dead streams
end
end, MatchingStreams),
{reply, ok, State};
handle_call({get_subscribers, Topic}, _From, State) when is_binary(Topic) ->
Subscribers = find_matching_subscribers(Topic, State),
{reply, {ok, Subscribers}, State};
handle_call({get_stream_topics, Stream}, _From, State) when is_pid(Stream) ->
Topics = maps:get(Stream, State#state.stream_subscriptions, []),
{reply, {ok, Topics}, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle stream process death - automatic cleanup.
handle_info({'DOWN', MonitorRef, process, StreamPid, _Reason}, State) ->
%% Remove monitor reference
Monitors = maps:remove(MonitorRef, State#state.monitors),
%% Get all topics this stream was subscribed to
Topics = maps:get(StreamPid, State#state.stream_subscriptions, []),
%% Remove stream from all topic subscriptions
NewSubscriptions = lists:foldl(fun(Topic, Acc) ->
Subscribers = maps:get(Topic, Acc, []),
NewSubscribers = lists:delete(StreamPid, Subscribers),
case NewSubscribers of
[] -> maps:remove(Topic, Acc);
_ -> maps:put(Topic, NewSubscribers, Acc)
end
end, State#state.subscriptions, Topics),
%% Remove stream from stream_subscriptions
NewStreamSubs = maps:remove(StreamPid, State#state.stream_subscriptions),
NewState = State#state{
subscriptions = NewSubscriptions,
stream_subscriptions = NewStreamSubs,
monitors = Monitors
},
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Find all streams subscribed to topics matching the given topic.
%% Supports exact matches and wildcard patterns (* and **).
-spec find_matching_subscribers(binary(), #state{}) -> [pid()].
find_matching_subscribers(Topic, State) ->
%% Get all subscription patterns
AllPatterns = maps:keys(State#state.subscriptions),
io:format("[PubSub ~p] Finding subscribers for topic: ~s~n", [self(), Topic]),
io:format("[PubSub ~p] All subscription patterns: ~p~n", [self(), AllPatterns]),
io:format("[PubSub ~p] Subscriptions map: ~p~n", [self(), State#state.subscriptions]),
%% Find patterns that match the topic
MatchingPatterns = lists:filter(fun(Pattern) ->
Matches = topic_matches(Pattern, Topic),
io:format("[PubSub ~p] Pattern ~s matches topic ~s: ~p~n", [self(), Pattern, Topic, Matches]),
Matches
end, AllPatterns),
io:format("[PubSub ~p] Matching patterns: ~p~n", [self(), MatchingPatterns]),
%% Collect all unique streams from matching patterns
AllStreams = lists:flatmap(fun(Pattern) ->
Streams = maps:get(Pattern, State#state.subscriptions, []),
io:format("[PubSub ~p] Pattern ~s has streams: ~p~n", [self(), Pattern, Streams]),
Streams
end, MatchingPatterns),
io:format("[PubSub ~p] All streams before dedup: ~p~n", [self(), AllStreams]),
%% Remove duplicates
Result = lists:usort(AllStreams),
io:format("[PubSub ~p] Final subscribers: ~p~n", [self(), Result]),
Result.
%% @doc Check if a topic pattern matches a concrete topic.
%% Supports:
%% - Exact match: "foo.bar" matches "foo.bar"
%% - Single-level wildcard: "foo.*.bar" matches "foo.xyz.bar"
%% - Multi-level wildcard: "foo.**.bar" matches "foo.x.y.z.bar"
-spec topic_matches(binary(), binary()) -> boolean().
topic_matches(Pattern, Topic) ->
%% Split into segments
PatternParts = binary:split(Pattern, <<".">>, [global]),
TopicParts = binary:split(Topic, <<".">>, [global]),
parts_match(PatternParts, TopicParts).
%% @doc Match pattern parts against topic parts.
-spec parts_match([binary()], [binary()]) -> boolean().
parts_match([], []) ->
true;
parts_match([<<"**">>], _) ->
%% ** matches any remaining segments
true;
parts_match([<<"**">> | PatternRest], TopicParts) ->
%% ** can match 0 or more segments
%% Try matching rest at every position
try_multi_wildcard(PatternRest, TopicParts);
parts_match([<<"*">> | PatternRest], [_TopicPart | TopicRest]) ->
%% * matches exactly one segment
parts_match(PatternRest, TopicRest);
parts_match([PatternPart | PatternRest], [TopicPart | TopicRest]) when PatternPart =:= TopicPart ->
%% Exact match
parts_match(PatternRest, TopicRest);
parts_match(_, _) ->
false.
%% @doc Try to match pattern after ** wildcard.
-spec try_multi_wildcard([binary()], [binary()]) -> boolean().
try_multi_wildcard(Pattern, Topic) ->
try_multi_wildcard(Pattern, Topic, 0).
-spec try_multi_wildcard([binary()], [binary()], non_neg_integer()) -> boolean().
try_multi_wildcard(Pattern, Topic, Skip) ->
case length(Topic) >= Skip of
true ->
TopicRest = lists:nthtail(Skip, Topic),
case parts_match(Pattern, TopicRest) of
true -> true;
false -> try_multi_wildcard(Pattern, Topic, Skip + 1)
end;
false ->
false
end.