Packages

macula

0.35.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
macula src macula_multi_relay.erl
Raw

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(),
%% Wrap callback with dedup filter
Self = self(),
DedupCallback = fun(Msg) ->
MsgId = extract_message_id(Msg),
case MsgId of
undefined -> Callback(Msg);
_ ->
case gen_server:call(Self, {check_dedup, MsgId}, 1000) of
new -> Callback(Msg);
duplicate -> ok
end
end
end,
Subs = maps:put(Ref, {Topic, Callback}, State#state.subscriptions),
%% Subscribe on all active connections
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_relay_client:subscribe(Pid, Topic, DedupCallback)
end, 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),
%% Unsubscribe from all connections
lists:foreach(fun(#conn{pid = Pid}) ->
%% We don't have the per-connection sub refs, so rely on
%% topic-level unsub. This works because relay_client tracks
%% by ref internally and we'd need a mapping. For now,
%% unsubscribe is best-effort (connection cleanup on terminate).
catch macula_relay_client:subscribe(Pid, Topic,
fun(_) -> ok end)
end, State#state.connections),
{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 = fun(Msg) ->
MsgId = extract_message_id(Msg),
case MsgId of
undefined -> Callback(Msg);
_ ->
case gen_server:call(Self, {check_dedup, MsgId}, 1000) of
new -> Callback(Msg);
duplicate -> ok
end
end
end,
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: 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])].