Packages
macula
0.6.6
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_server.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Pub/Sub GenServer - manages subscriptions and message delivery.
%%% Ties together registry, cache, discovery, and delivery layers.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_pubsub_server).
-behaviour(gen_server).
%% API
-export([
start_link/0,
start_link/1,
stop/1,
subscribe/4,
unsubscribe/3,
publish/2,
list_subscriptions/1,
list_patterns/1,
subscription_count/1,
cache_stats/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
%% Types
-record(state, {
registry :: macula_pubsub_registry:registry(),
cache :: macula_pubsub_cache:cache(),
cache_ttl :: pos_integer(),
discovery_fun :: macula_pubsub_discovery:dht_lookup_fun() | undefined,
send_fun :: macula_pubsub_delivery:send_fun() | undefined
}).
-type options() :: #{
cache_size => pos_integer(),
cache_ttl => pos_integer(),
discovery_fun => macula_pubsub_discovery:dht_lookup_fun(),
send_fun => macula_pubsub_delivery:send_fun()
}.
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start server with default options.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start server with options.
-spec start_link(options()) -> {ok, pid()} | {error, term()}.
start_link(Options) ->
gen_server:start_link(?MODULE, Options, []).
%% @doc Stop server.
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Subscribe to a pattern.
-spec subscribe(pid(), binary(), binary(), pid()) -> ok.
subscribe(Pid, SubscriberId, Pattern, Callback) ->
gen_server:call(Pid, {subscribe, SubscriberId, Pattern, Callback}).
%% @doc Unsubscribe from a pattern.
-spec unsubscribe(pid(), binary(), binary()) -> ok.
unsubscribe(Pid, SubscriberId, Pattern) ->
gen_server:call(Pid, {unsubscribe, SubscriberId, Pattern}).
%% @doc Publish message to all matching subscribers.
-spec publish(pid(), macula_pubsub_delivery:message()) -> ok.
publish(Pid, Message) ->
gen_server:cast(Pid, {publish, Message}).
%% @doc List all subscriptions.
-spec list_subscriptions(pid()) -> [macula_pubsub_registry:subscription()].
list_subscriptions(Pid) ->
gen_server:call(Pid, list_subscriptions).
%% @doc List all unique patterns.
-spec list_patterns(pid()) -> [binary()].
list_patterns(Pid) ->
gen_server:call(Pid, list_patterns).
%% @doc Get subscription count.
-spec subscription_count(pid()) -> non_neg_integer().
subscription_count(Pid) ->
gen_server:call(Pid, subscription_count).
%% @doc Get cache statistics.
-spec cache_stats(pid()) -> #{size := non_neg_integer(), max_size := pos_integer()}.
cache_stats(Pid) ->
gen_server:call(Pid, cache_stats).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init(Options) ->
%% Extract options
CacheSize = maps:get(cache_size, Options, 1000),
CacheTTL = maps:get(cache_ttl, Options, 300),
DiscoveryFun = maps:get(discovery_fun, Options, undefined),
SendFun = maps:get(send_fun, Options, undefined),
%% Initialize state
State = #state{
registry = macula_pubsub_registry:new(),
cache = macula_pubsub_cache:new(CacheSize),
cache_ttl = CacheTTL,
discovery_fun = DiscoveryFun,
send_fun = SendFun
},
{ok, State}.
%% @private
handle_call({subscribe, SubscriberId, Pattern, Callback}, _From, State) ->
#state{registry = Registry} = State,
NewRegistry = macula_pubsub_registry:subscribe(Registry, SubscriberId, Pattern, Callback),
{reply, ok, State#state{registry = NewRegistry}};
handle_call({unsubscribe, SubscriberId, Pattern}, _From, State) ->
#state{registry = Registry} = State,
NewRegistry = macula_pubsub_registry:unsubscribe(Registry, SubscriberId, Pattern),
{reply, ok, State#state{registry = NewRegistry}};
handle_call(list_subscriptions, _From, State) ->
#state{registry = Registry} = State,
%% Get all subscriptions from registry
Subscriptions = get_all_subscriptions(Registry),
{reply, Subscriptions, State};
handle_call(list_patterns, _From, State) ->
#state{registry = Registry} = State,
Patterns = macula_pubsub_registry:list_patterns(Registry),
{reply, Patterns, State};
handle_call(subscription_count, _From, State) ->
#state{registry = Registry} = State,
Count = macula_pubsub_registry:size(Registry),
{reply, Count, State};
handle_call(cache_stats, _From, State) ->
#state{cache = Cache} = State,
Stats = #{
size => macula_pubsub_cache:size(Cache),
max_size => macula_pubsub_cache:max_size(Cache)
},
{reply, Stats, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @private
handle_cast({publish, Message}, State) ->
#state{
registry = Registry,
discovery_fun = DiscoveryFun,
send_fun = SendFun
} = State,
%% Create discovery function (use cache if available)
DiscoveryFn = get_discovery_fn(DiscoveryFun),
%% Create send function (default: no-op)
SendFn = get_send_fn(SendFun),
%% Publish to local and remote
_Results = macula_pubsub_delivery:publish(Message, Registry, DiscoveryFn, SendFn),
{noreply, State};
handle_cast(_Request, State) ->
{noreply, State}.
%% @private
handle_info(_Info, State) ->
{noreply, State}.
%% @private
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Get all subscriptions from registry.
%% NOTE: This extracts subscriptions from the registry internal structure.
-spec get_all_subscriptions(macula_pubsub_registry:registry()) ->
[macula_pubsub_registry:subscription()].
get_all_subscriptions(#{subscriptions := Subs}) ->
Subs.
%% @doc Get discovery function or default.
get_discovery_fn(undefined) -> fun(_Pattern) -> {ok, []} end;
get_discovery_fn(Fun) -> Fun.
%% @doc Get send function or default.
get_send_fn(undefined) -> fun(_Msg, _Addr) -> ok end;
get_send_fn(Fn) -> Fn.