Packages
macula
1.4.4
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_mesh_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_mesh_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, call_any/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
relay_discovery_subscribed = false :: boolean(), %% subscribed to _mesh.relay.up/down
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_table :: ets:tid(), %% ETS set for O(1) concurrent dedup
dedup_queue :: queue:queue(binary()) %% eviction order tracking
}).
%%====================================================================
%% API — same signatures as macula_mesh_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).
%% @doc Try an RPC call on each connected relay in sequence.
%% Returns the first successful result or the last error.
%% Useful for cross-relay operations where the target may be on a different relay.
-spec call_any(pid(), binary(), map(), timeout()) -> {ok, term()} | {error, term()}.
call_any(Pid, Procedure, Args, Timeout) ->
gen_server:call(Pid, {rpc_call_any, Procedure, Args, Timeout}, Timeout + 5000).
-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,
maybe_start_discovery(Relays, Opts),
TargetCount = min(
maps:get(connections, Opts, ?DEFAULT_CONNECTIONS),
length(Relays)
),
DedupTable = ets:new(multi_relay_dedup, [set, public, {read_concurrency, true}]),
State = #state{
opts = Opts,
relays = Relays,
connections = [],
target_count = TargetCount,
subscriptions = #{},
procedures = #{},
dedup_table = DedupTable,
dedup_queue = queue:new()
},
self() ! spawn_connections,
erlang:send_after(?HEALTH_CHECK_MS, self(), health_check),
{ok, State}.
%% Start relay discovery if the caller provides geographic coordinates.
%% Discovery runs as a standalone gen_server — relay_client uses it for failover.
maybe_start_discovery(Seeds, #{site := #{lat := Lat, lng := Lng}}) when is_number(Lat), is_number(Lng) ->
case whereis(macula_relay_discovery) of
undefined ->
macula_relay_discovery:start_link(#{seeds => Seeds, lat => Lat, lng => Lng}),
?LOG_INFO("[multi_relay] Started relay discovery (lat: ~.2f, lng: ~.2f)", [Lat, Lng]);
_ ->
ok
end;
maybe_start_discovery(_, _) ->
ok.
%%====================================================================
%% Subscribe / Unsubscribe — applied to ALL connections
%%====================================================================
handle_call({subscribe, Topic, Callback}, _From, State) ->
Ref = make_ref(),
DedupCallback = make_dedup_callback(State#state.dedup_table, 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}) ->
case macula_mesh_client:advertise(Pid, Procedure, Handler) of
{ok, _} -> ok;
{error, Reason} ->
?LOG_WARNING("[multi_relay] Advertise ~s on ~p failed: ~p",
[Procedure, Pid, Reason])
end
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}) ->
case macula_mesh_client:unadvertise(Pid, Procedure) of
ok -> ok;
{error, Reason} ->
?LOG_WARNING("[multi_relay] Unadvertise ~s on ~p failed: ~p",
[Procedure, Pid, Reason])
end
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_mesh_client:call(Pid, Procedure, Args, Timeout),
{reply, Result, State}
end;
handle_call({rpc_call_any, Procedure, Args, Timeout}, _From, State) ->
AliveConns = [C || C <- State#state.connections, is_process_alive(C#conn.pid)],
Result = try_call_each(AliveConns, Procedure, Args, Timeout),
{reply, Result, State};
%%====================================================================
%% Status
%%====================================================================
handle_call(get_status, _From, State) ->
Conns = [#{
relay => C#conn.relay_url,
role => atom_to_binary(C#conn.role),
alive => is_process_alive(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 => ets:info(State#state.dedup_table, size)
},
{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_mesh_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),
State3 = maybe_subscribe_relay_events(State2),
{noreply, State3};
handle_info(health_check, State) ->
State2 = ensure_connections(State),
State3 = cleanup_dedup(State2),
erlang:send_after(?HEALTH_CHECK_MS, self(), health_check),
{noreply, State3};
%% 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;
%% Relay discovery events — forward to macula_relay_discovery
handle_info({relay_event, <<"_mesh.relay.up">>, Payload}, State) ->
forward_relay_event(relay_up, Payload),
{noreply, State};
handle_info({relay_event, <<"_mesh.relay.down">>, Payload}, State) ->
forward_relay_event(relay_down, Payload),
{noreply, State};
handle_info(_Info, State) ->
{noreply, State}.
%%====================================================================
%% Relay discovery integration
%%====================================================================
%% Subscribe to _mesh.relay.up/down once we have connections and discovery is running.
maybe_subscribe_relay_events(#state{relay_discovery_subscribed = true} = State) ->
State;
maybe_subscribe_relay_events(#state{connections = []} = State) ->
State;
maybe_subscribe_relay_events(State) ->
case whereis(macula_relay_discovery) of
undefined ->
State;
_ ->
Self = self(),
Primary = get_primary(State),
subscribe_relay_topic(Primary, <<"_mesh.relay.up">>, Self),
subscribe_relay_topic(Primary, <<"_mesh.relay.down">>, Self),
?LOG_INFO("[multi_relay] Subscribed to relay discovery events"),
State#state{relay_discovery_subscribed = true}
end.
subscribe_relay_topic(undefined, _, _) -> ok;
subscribe_relay_topic(#conn{pid = Pid}, Topic, Self) ->
%% Non-blocking — the mesh client may still be connecting (QUIC NIF).
%% A synchronous call here causes a timeout crash if connect is slow.
spawn(fun() ->
try macula_mesh_client:subscribe(Pid, Topic, fun(Msg) ->
Self ! {relay_event, Topic, Msg}
end)
catch _:_ -> ok
end
end).
get_primary(#state{connections = Conns}) ->
lists:keyfind(primary, #conn.role, Conns).
forward_relay_event(relay_up, #{<<"hostname">> := Hostname} = Info) ->
Lat = maps:get(<<"lat">>, Info, 0.0),
Lng = maps:get(<<"lng">>, Info, 0.0),
macula_relay_discovery ! {relay_up, Hostname, #{lat => Lat, lng => Lng}};
forward_relay_event(relay_down, #{<<"hostname">> := Hostname}) ->
macula_relay_discovery ! {relay_down, Hostname};
forward_relay_event(_, _) ->
ok.
terminate(_Reason, #state{connections = Conns}) ->
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_mesh_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,
dedup_table = DedupTable}) ->
%% Build opts for this specific relay connection
ClientOpts = Opts#{
relays => [RelayUrl],
url => RelayUrl
},
case macula_mesh_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, DedupTable),
%% 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, DedupTable) ->
maps:foreach(fun(_Ref, {Topic, Callback}) ->
DedupCallback = make_dedup_callback(DedupTable, Callback),
catch macula_mesh_client:subscribe(Pid, Topic, DedupCallback)
end, Subs).
replay_procedures(Pid, Procs) ->
maps:foreach(fun(Procedure, Handler) ->
catch macula_mesh_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;
[] ->
case Conns of
[First | _] -> First;
[] -> undefined
end
end.
%% Try RPC on each connection in sequence. Return first success or last error.
try_call_each([], _Procedure, _Args, _Timeout) ->
{error, no_connection};
try_call_each([#conn{pid = Pid}], Procedure, Args, Timeout) ->
macula_mesh_client:call(Pid, Procedure, Args, Timeout);
try_call_each([#conn{pid = Pid} | Rest], Procedure, Args, Timeout) ->
case macula_mesh_client:call(Pid, Procedure, Args, Timeout) of
{ok, _} = Success -> Success;
{error, _} -> try_call_each(Rest, Procedure, Args, Timeout)
end.
%%====================================================================
%% Internal: subscription helpers
%%====================================================================
subscribe_on_all(Topic, Callback, Connections) ->
lists:foreach(fun(#conn{pid = Pid}) ->
catch macula_mesh_client:subscribe(Pid, Topic, Callback)
end, Connections).
make_dedup_callback(DedupTable, Callback) ->
fun(Msg) ->
maybe_invoke_deduped(DedupTable, Callback, Msg)
end.
maybe_invoke_deduped(DedupTable, Callback, Msg) ->
case extract_message_id(Msg) of
undefined -> Callback(Msg);
MsgId -> invoke_if_new(DedupTable, Callback, Msg, MsgId)
end.
%% O(1) concurrent dedup via public ETS — no gen_server bottleneck.
%% insert_new is atomic: returns true only for the first inserter.
invoke_if_new(DedupTable, Callback, Msg, MsgId) ->
case ets:insert_new(DedupTable, {MsgId, true}) of
true -> Callback(Msg);
false -> ok
end.
%% Note: ETS cleanup happens periodically in health_check, not per-message.
%% The queue tracking for eviction order is maintained in the gen_server
%% state but doesn't block the hot path (insert_new is lock-free).
%%====================================================================
%% Internal: dedup ETS cleanup
%%====================================================================
%% Periodic cleanup: keep ETS table bounded by evicting oldest entries.
%% Called from health_check timer. Uses the queue to track insertion order.
cleanup_dedup(#state{dedup_table = Tab, dedup_queue = Q} = State) ->
Size = ets:info(Tab, size),
case Size > ?DEDUP_SIZE of
false -> State;
true ->
{Q2, _Evicted} = evict_oldest(Tab, Q, Size - ?DEDUP_SIZE),
State#state{dedup_queue = Q2}
end.
evict_oldest(_Tab, Q, 0) -> {Q, 0};
evict_oldest(Tab, Q, N) ->
case queue:out(Q) of
{{value, MsgId}, Q2} ->
ets:delete(Tab, MsgId),
evict_oldest(Tab, Q2, N - 1);
{empty, Q} ->
{Q, 0}
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])].