Packages
macula
0.35.2
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_multi_relay.erl
%%%-------------------------------------------------------------------
%%% @doc Multi-relay client — maintains N concurrent relay connections.
%%%
%%% Wraps multiple macula_relay_client instances for node multi-homing.
%%% Subscribes and advertises on ALL connections. Publishes via PRIMARY.
%%% Deduplicates incoming messages by message_id (ring buffer).
%%%
%%% On primary failure: promotes secondary, connects a new secondary.
%%% Same API as macula_relay_client for drop-in replacement.
%%%
%%% Usage:
%%% ```
%%% Opts = #{relays => [R1, R2, R3, R4, R5],
%%% connections => 2, %% default 2
%%% realm => <<"io.macula">>,
%%% identity => <<"my-node">>},
%%% {ok, Pid} = macula_multi_relay:start_link(Opts).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_multi_relay).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([start_link/1, stop/1]).
-export([subscribe/3, unsubscribe/2, publish/3]).
-export([advertise/3, unadvertise/2, call/4]).
-export([get_status/1]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
%% Test exports
-ifdef(TEST).
-export([extract_message_id/1, assign_roles/1, shuffle/1]).
-endif.
-define(DEDUP_SIZE, 2048).
-define(DEFAULT_CONNECTIONS, 2).
-define(HEALTH_CHECK_MS, 15000).
-record(conn, {
pid :: pid(),
relay_url :: binary(),
monitor :: reference(),
role :: primary | secondary
}).
-record(state, {
opts :: map(), %% original start opts (for spawning new connections)
relays :: [binary()], %% all configured relay URLs
connections :: [#conn{}], %% active connections
target_count :: pos_integer(), %% how many connections to maintain
subscriptions :: #{reference() => {binary(), fun()}}, %% ref => {topic, callback}
procedures :: #{binary() => fun()}, %% procedure => handler
dedup :: queue:queue(binary()), %% ring buffer of recent message_ids
dedup_set :: sets:set(binary()) %% fast lookup for dedup
}).
%%====================================================================
%% API — same signatures as macula_relay_client
%%====================================================================
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
stop(Pid) ->
gen_server:stop(Pid).
-spec subscribe(pid(), binary(), fun((map()) -> ok)) -> {ok, reference()}.
subscribe(Pid, Topic, Callback) ->
gen_server:call(Pid, {subscribe, Topic, Callback}).
-spec unsubscribe(pid(), reference()) -> ok.
unsubscribe(Pid, Ref) ->
gen_server:call(Pid, {unsubscribe, Ref}).
-spec publish(pid(), binary(), binary() | map()) -> ok.
publish(Pid, Topic, Payload) ->
gen_server:cast(Pid, {publish, Topic, Payload}).
-spec advertise(pid(), binary(), fun((map()) -> {ok, term()} | {error, term()})) -> {ok, reference()}.
advertise(Pid, Procedure, Handler) ->
gen_server:call(Pid, {advertise, Procedure, Handler}).
-spec unadvertise(pid(), binary()) -> ok.
unadvertise(Pid, Procedure) ->
gen_server:call(Pid, {unadvertise, Procedure}).
-spec call(pid(), binary(), map(), timeout()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args, Timeout) ->
gen_server:call(Pid, {rpc_call, Procedure, Args, Timeout}, Timeout + 1000).
-spec get_status(pid()) -> {ok, map()}.
get_status(Pid) ->
gen_server:call(Pid, get_status).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init(Opts) ->
Relays = case maps:find(relays, Opts) of
{ok, List} when is_list(List), length(List) > 0 -> List;
_ -> [maps:get(url, Opts, <<"https://localhost:4433">>)]
end,
TargetCount = min(
maps:get(connections, Opts, ?DEFAULT_CONNECTIONS),
length(Relays)
),
State = #state{
opts = Opts,
relays = Relays,
connections = [],
target_count = TargetCount,
subscriptions = #{},
procedures = #{},
dedup = queue:new(),
dedup_set = sets:new([{version, 2}])
},
%% Start connections asynchronously
self() ! spawn_connections,
%% Periodic health check
erlang:send_after(?HEALTH_CHECK_MS, self(), health_check),
{ok, State}.
%%====================================================================
%% Subscribe / Unsubscribe — applied to ALL connections
%%====================================================================
handle_call({subscribe, Topic, Callback}, _From, State) ->
Ref = make_ref(),
DedupCallback = make_dedup_callback(self(), Callback),
Subs = maps:put(Ref, {Topic, Callback}, State#state.subscriptions),
subscribe_on_all(Topic, DedupCallback, State#state.connections),
{reply, {ok, Ref}, State#state{subscriptions = Subs}};
handle_call({unsubscribe, Ref}, _From, State) ->
case maps:get(Ref, State#state.subscriptions, undefined) of
undefined ->
{reply, {error, not_found}, State};
{_Topic, _Callback} ->
Subs = maps:remove(Ref, State#state.subscriptions),
%% Note: per-connection sub refs aren't tracked in multi_relay.
%% Subscriptions are cleaned up when connections terminate.
{reply, ok, State#state{subscriptions = Subs}}
end;
%%====================================================================
%% Advertise / Unadvertise — applied to ALL connections
%%====================================================================
handle_call({advertise, Procedure, Handler}, _From, State) ->
Procs = maps:put(Procedure, Handler, State#state.procedures),
%% Advertise on all active connections
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_relay_client:advertise(Pid, Procedure, Handler)
end, State#state.connections),
{reply, {ok, make_ref()}, State#state{procedures = Procs}};
handle_call({unadvertise, Procedure}, _From, State) ->
Procs = maps:remove(Procedure, State#state.procedures),
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_relay_client:unadvertise(Pid, Procedure)
end, State#state.connections),
{reply, ok, State#state{procedures = Procs}};
%%====================================================================
%% RPC Call — via primary connection
%%====================================================================
handle_call({rpc_call, Procedure, Args, Timeout}, _From, State) ->
case find_primary(State#state.connections) of
undefined ->
{reply, {error, no_connection}, State};
#conn{pid = Pid} ->
Result = macula_relay_client:call(Pid, Procedure, Args, Timeout),
{reply, Result, State}
end;
%%====================================================================
%% Dedup check (called from callback wrapper)
%%====================================================================
handle_call({check_dedup, MsgId}, _From, State) ->
case sets:is_element(MsgId, State#state.dedup_set) of
true ->
{reply, duplicate, State};
false ->
{Ring, Set} = track_message(MsgId, State#state.dedup, State#state.dedup_set),
{reply, new, State#state{dedup = Ring, dedup_set = Set}}
end;
%%====================================================================
%% Status
%%====================================================================
handle_call(get_status, _From, State) ->
Conns = [#{
relay => C#conn.relay_url,
role => C#conn.role,
pid => C#conn.pid
} || C <- State#state.connections],
Status = #{
connections => Conns,
target_count => State#state.target_count,
active_count => length(State#state.connections),
subscriptions => maps:size(State#state.subscriptions),
procedures => maps:size(State#state.procedures),
dedup_size => queue:len(State#state.dedup)
},
{reply, {ok, Status}, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
%%====================================================================
%% Publish — via primary only (avoid duplicate publishes)
%%====================================================================
handle_cast({publish, Topic, Payload}, State) ->
case find_primary(State#state.connections) of
undefined -> ok;
#conn{pid = Pid} ->
macula_relay_client:publish(Pid, Topic, Payload)
end,
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%%====================================================================
%% Connection lifecycle
%%====================================================================
handle_info(spawn_connections, State) ->
State2 = ensure_connections(State),
{noreply, State2};
handle_info(health_check, State) ->
State2 = ensure_connections(State),
erlang:send_after(?HEALTH_CHECK_MS, self(), health_check),
{noreply, State2};
%% A relay_client process died — remove from connections, spawn replacement
handle_info({'DOWN', MonRef, process, _Pid, Reason}, State) ->
case lists:keyfind(MonRef, #conn.monitor, State#state.connections) of
false ->
{noreply, State};
#conn{relay_url = Url, role = Role} ->
?LOG_WARNING("[multi_relay] Connection to ~s (~p) down: ~p",
[Url, Role, Reason]),
Conns = lists:keydelete(MonRef, #conn.monitor, State#state.connections),
State2 = State#state{connections = Conns},
%% If primary died, promote first secondary
State3 = case Role of
primary -> promote_secondary(State2);
secondary -> State2
end,
%% Schedule replacement connection
erlang:send_after(1000, self(), spawn_connections),
{noreply, State3}
end;
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, #state{connections = Conns}) ->
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_relay_client:stop(Pid)
end, Conns),
ok.
%%====================================================================
%% Internal: connection management
%%====================================================================
ensure_connections(#state{connections = Conns, target_count = Target} = State)
when length(Conns) >= Target ->
State;
ensure_connections(#state{connections = Conns, target_count = Target,
relays = Relays, opts = Opts} = State) ->
Needed = Target - length(Conns),
UsedUrls = [C#conn.relay_url || C <- Conns],
%% Pick relays not already connected to
Available = [R || R <- Relays, not lists:member(R, UsedUrls)],
%% If all relays are used, allow duplicates (fewer relays than connections)
Candidates = case Available of
[] -> Relays;
_ -> Available
end,
%% Shuffle for load distribution
Shuffled = shuffle(Candidates),
ToConnect = lists:sublist(Shuffled, Needed),
NewConns = lists:foldl(fun(RelayUrl, Acc) ->
case spawn_relay_client(RelayUrl, Opts, State) of
{ok, Conn} -> [Conn | Acc];
{error, _} -> Acc
end
end, [], ToConnect),
AllConns = Conns ++ NewConns,
%% Assign roles: first is primary, rest are secondary
Tagged = assign_roles(AllConns),
State#state{connections = Tagged}.
spawn_relay_client(RelayUrl, Opts, #state{subscriptions = Subs, procedures = Procs}) ->
Self = self(),
%% Build opts for this specific relay connection
ClientOpts = Opts#{
relays => [RelayUrl],
url => RelayUrl
},
case macula_relay_client:start_link(ClientOpts) of
{ok, Pid} ->
MonRef = monitor(process, Pid),
?LOG_INFO("[multi_relay] Connected to ~s (pid ~p)", [RelayUrl, Pid]),
%% Replay subscriptions on this new connection
replay_subscriptions(Pid, Subs, Self),
%% Replay procedures on this new connection
replay_procedures(Pid, Procs),
{ok, #conn{pid = Pid, relay_url = RelayUrl, monitor = MonRef, role = secondary}};
{error, Reason} ->
?LOG_WARNING("[multi_relay] Failed to start client for ~s: ~p",
[RelayUrl, Reason]),
{error, Reason}
end.
replay_subscriptions(Pid, Subs, Self) ->
maps:foreach(fun(_Ref, {Topic, Callback}) ->
DedupCallback = make_dedup_callback(Self, Callback),
catch macula_relay_client:subscribe(Pid, Topic, DedupCallback)
end, Subs).
replay_procedures(Pid, Procs) ->
maps:foreach(fun(Procedure, Handler) ->
catch macula_relay_client:advertise(Pid, Procedure, Handler)
end, Procs).
assign_roles([]) -> [];
assign_roles([First | Rest]) ->
[First#conn{role = primary} | [C#conn{role = secondary} || C <- Rest]].
promote_secondary(#state{connections = []} = State) ->
?LOG_WARNING("[multi_relay] No connections remaining — cannot promote"),
State;
promote_secondary(#state{connections = [First | Rest]} = State) ->
?LOG_INFO("[multi_relay] Promoting ~s to primary", [First#conn.relay_url]),
State#state{connections = [First#conn{role = primary} | Rest]}.
find_primary(Conns) ->
case [C || #conn{role = primary} = C <- Conns] of
[Primary | _] -> Primary;
[] ->
%% No primary tagged — use first available
case Conns of
[First | _] -> First;
[] -> undefined
end
end.
%%====================================================================
%% Internal: subscription helpers
%%====================================================================
subscribe_on_all(Topic, Callback, Connections) ->
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_relay_client:subscribe(Pid, Topic, Callback)
end, Connections).
make_dedup_callback(MultiRelayPid, Callback) ->
fun(Msg) ->
maybe_invoke_deduped(MultiRelayPid, Callback, Msg)
end.
maybe_invoke_deduped(MultiRelayPid, Callback, Msg) ->
case extract_message_id(Msg) of
undefined -> Callback(Msg);
MsgId -> invoke_if_new(MultiRelayPid, Callback, Msg, MsgId)
end.
invoke_if_new(MultiRelayPid, Callback, Msg, MsgId) ->
case gen_server:call(MultiRelayPid, {check_dedup, MsgId}, 1000) of
new -> Callback(Msg);
duplicate -> ok
end.
%%====================================================================
%% Internal: dedup ring buffer
%%====================================================================
track_message(MsgId, Ring, Set) ->
Ring1 = queue:in(MsgId, Ring),
Set1 = sets:add_element(MsgId, Set),
case queue:len(Ring1) > ?DEDUP_SIZE of
true ->
{{value, OldId}, Ring2} = queue:out(Ring1),
Set2 = sets:del_element(OldId, Set1),
{Ring2, Set2};
false ->
{Ring1, Set1}
end.
extract_message_id(#{<<"message_id">> := MsgId}) when is_binary(MsgId) -> MsgId;
extract_message_id(#{message_id := MsgId}) when is_binary(MsgId) -> MsgId;
extract_message_id(#{payload := Payload}) when is_map(Payload) ->
case maps:find(<<"message_id">>, Payload) of
{ok, MsgId} when is_binary(MsgId) -> MsgId;
_ -> undefined
end;
extract_message_id(_) -> undefined.
%%====================================================================
%% Internal: utils
%%====================================================================
shuffle(List) ->
[X || {_, X} <- lists:sort([{rand:uniform(), E} || E <- List])].