Packages
macula
0.10.1
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_gateway_system/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]).
%% Stream handle can be either a pid (process) or reference (QUIC stream)
-type stream_handle() :: pid() | reference().
-record(state, {
opts :: map(),
subscriptions :: #{binary() => [stream_handle()]}, % topic => [stream_handles]
stream_subscriptions :: #{stream_handle() => [binary()]}, % stream_handle => [topics]
monitors :: #{reference() => {stream_handle(), binary()}} % monitor_ref => {stream_handle, 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).
%% Async (cast) to prevent blocking callers when PubSub is busy.
-spec subscribe(pid(), pid() | reference(), binary()) -> ok.
subscribe(Pid, Stream, Topic) ->
gen_server:cast(Pid, {subscribe, Stream, Topic}).
%% @doc Unsubscribe a stream from a topic.
%% Async (cast) to prevent blocking callers when PubSub is busy.
-spec unsubscribe(pid(), pid() | reference(), binary()) -> ok.
unsubscribe(Pid, Stream, Topic) ->
gen_server:cast(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"),
%% DEBUG: Log what's in Opts to verify node_id/url are present
io:format("[PubSub DEBUG] Opts map keys: ~p~n", [maps:keys(Opts)]),
io:format("[PubSub DEBUG] node_id in opts: ~p~n", [maps:get(node_id, Opts, not_found)]),
io:format("[PubSub DEBUG] url in opts: ~p~n", [maps:get(url, Opts, not_found)]),
State = #state{
opts = Opts,
subscriptions = #{},
stream_subscriptions = #{},
monitors = #{}
},
%% Schedule re-advertisement of all subscriptions after DHT routing table is populated
%% Initial subscriptions happen before bootstrap connects, so this ensures they propagate
erlang:send_after(5000, self(), readvertise_all_subscriptions),
io:format("[PubSub] Pub/sub handler initialized~n"),
{ok, State}.
%% Subscribe is async (cast) to prevent blocking callers when PubSub is busy
handle_cast({subscribe, Stream, Topic}, State) when (is_pid(Stream) orelse is_reference(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])]),
%% Advertise subscription in DHT for cross-gateway discovery
advertise_subscription_in_dht(Topic, State),
%% Add to stream → topics mapping
NewTopics = [Topic | CurrentTopics],
NewStreamSubs = maps:put(Stream, NewTopics, State#state.stream_subscriptions),
%% Monitor stream if first subscription (only for pids - can't monitor references)
NewMonitors = case {length(CurrentTopics), is_pid(Stream)} of
{0, true} ->
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,
{noreply, NewState};
%% Unsubscribe is async (cast) to prevent blocking callers when PubSub is busy
handle_cast({unsubscribe, Stream, Topic}, State) when (is_pid(Stream) orelse is_reference(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
},
{noreply, NewState};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_call({publish, Topic, Payload}, _From, State) when is_binary(Topic) ->
%% Find all LOCAL matching subscribers (exact + wildcard patterns)
LocalStreams = find_matching_subscribers(Topic, State),
io:format("[PubSub ~p] Publishing to topic ~s: found ~p local subscribers~n",
[self(), Topic, length(LocalStreams)]),
%% Get node_id for distribute call
LocalNodeId = maps:get(node_id, State#state.opts, <<"unknown">>),
%% Create PUBLISH message map as expected by pubsub_router
PubMsg = #{
<<"topic">> => Topic,
<<"payload">> => Payload
},
%% Use the pubsub_router to distribute to both local and remote subscribers
%% This properly uses DHT routing + mesh connections + pubsub_route envelopes
case {whereis(macula_gateway_mesh), whereis(macula_gateway_client_manager)} of
{undefined, _} ->
io:format("[PubSub ~p] WARNING: macula_gateway_mesh not running, cannot route remotely~n", [self()]),
%% Fallback: at least deliver locally
lists:foreach(fun(Stream) ->
case erlang:is_process_alive(Stream) of
true ->
io:format("[PubSub ~p] Sending to local stream ~p~n", [self(), Stream]),
Stream ! {publish, Topic, Payload};
false ->
ok
end
end, LocalStreams);
{_, undefined} ->
io:format("[PubSub ~p] WARNING: macula_gateway_client_manager not running, cannot route remotely~n", [self()]),
%% Fallback: at least deliver locally
lists:foreach(fun(Stream) ->
case erlang:is_process_alive(Stream) of
true ->
io:format("[PubSub ~p] Sending to local stream ~p~n", [self(), Stream]),
Stream ! {publish, Topic, Payload};
false ->
ok
end
end, LocalStreams);
{MeshPid, ClientsPid} ->
%% Use pubsub_router for proper DHT-based routing
io:format("[PubSub ~p] Using pubsub_router to distribute (mesh=~p, clients=~p)~n",
[self(), MeshPid, ClientsPid]),
macula_gateway_pubsub_router:distribute(
LocalStreams,
PubMsg,
LocalNodeId,
MeshPid,
ClientsPid
)
end,
{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) orelse is_reference(Stream) ->
Topics = maps:get(Stream, State#state.stream_subscriptions, []),
{reply, {ok, Topics}, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, 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};
%% @doc Re-advertise all existing subscriptions in the DHT.
%% Called after bootstrap connection to ensure subscriptions propagate.
handle_info(readvertise_all_subscriptions, State) ->
Topics = maps:keys(State#state.subscriptions),
io:format("[PubSub] Re-advertising ~p subscription(s) after DHT routing table populated~n",
[length(Topics)]),
lists:foreach(fun(Topic) ->
advertise_subscription_in_dht(Topic, State)
end, Topics),
{noreply, State};
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.
%% @private
%% @doc Advertise a subscription in the DHT so other gateways can discover it.
%% This enables cross-gateway pub/sub routing.
%% Also invalidates subscriber cache to ensure fresh lookups.
-spec advertise_subscription_in_dht(binary(), #state{}) -> ok.
advertise_subscription_in_dht(Topic, State) ->
%% Invalidate subscriber cache - subscription change means cached subscribers are stale
invalidate_subscriber_cache(Topic),
%% Get gateway node_id and endpoint from opts
NodeId = maps:get(node_id, State#state.opts, <<"unknown">>),
Url = maps:get(url, State#state.opts, <<"unknown">>),
%% DEBUG: Log what we got from State#state.opts
io:format("[GatewayPubSub DEBUG] Advertising subscription for topic ~s~n", [Topic]),
io:format("[GatewayPubSub DEBUG] NodeId from opts: ~p~n", [NodeId]),
io:format("[GatewayPubSub DEBUG] Url from opts: ~p~n", [Url]),
%% Create DHT key by hashing the topic
TopicKey = crypto:hash(sha256, Topic),
%% Create subscriber value with endpoint info for routing messages back
SubscriberValue = #{
node_id => NodeId,
endpoint => Url,
ttl => 300 % 5 minutes TTL
},
%% Store subscription in DHT (will propagate to k closest nodes or bootstrap)
io:format("[GatewayPubSub] Advertising subscription to topic ~s in DHT~n", [Topic]),
try
case whereis(macula_routing_server) of
undefined ->
io:format("[GatewayPubSub] Routing server not running, cannot advertise subscription~n");
RoutingServerPid ->
case macula_routing_server:store(RoutingServerPid, TopicKey, SubscriberValue) of
ok ->
io:format("[GatewayPubSub] Successfully stored subscription for ~s in DHT~n", [Topic]);
{error, StoreError} ->
io:format("[GatewayPubSub] Failed to store subscription ~s: ~p~n",
[Topic, StoreError])
end
end
catch
_:DhtError ->
io:format("[GatewayPubSub] Failed to advertise subscription ~s in DHT: ~p (continuing)~n",
[Topic, DhtError])
end,
ok.
%% @private
%% @doc Invalidate subscriber cache for a topic.
%% Called when subscription changes to ensure fresh DHT lookups.
-spec invalidate_subscriber_cache(binary()) -> ok.
invalidate_subscriber_cache(Topic) ->
case whereis(macula_subscriber_cache) of
undefined ->
%% Cache not running yet - that's fine
ok;
_Pid ->
macula_subscriber_cache:invalidate(Topic)
end.