Packages

macula

0.45.1
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_relay_client.erl
Raw

src/macula_relay_client.erl

%%%-------------------------------------------------------------------
%%% @doc Macula Relay Client — connects to relay servers via QUIC.
%%%
%%% Single gen_server managing ONE persistent QUIC connection with
%%% automatic failover across multiple relays. On disconnect, cycles
%%% to the next relay in the list with exponential backoff + jitter.
%%% Replays all subscriptions and procedure registrations on reconnect.
%%%
%%% Usage:
%%%
%%% ```
%%% Opts = #{relays => [Relay0, Relay1, Relay2]},
%%% {ok, Client} = macula_relay_client:start_link(Opts).
%%% '''
%%%
%%% Backward compatible with single relay via url key.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_relay_client).
-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([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
%% Test exports
-ifdef(TEST).
-export([backoff_ms/1, parse_url/1]).
-endif.
-define(RECONNECT_BASE_MS, 1000).
-define(RECONNECT_MAX_MS, 30000).
-define(RECONNECT_JITTER, 0.3).
-define(CALL_TIMEOUT, 5000).
-define(PING_INTERVAL_MS, 30000).
-define(PING_TOPIC, <<"_mesh.relay.ping">>).
-record(state, {
relays :: [binary()], %% all configured relay URLs
relay_index :: non_neg_integer(), %% current relay (0-based)
url :: binary(), %% current relay URL
host :: string(),
port :: integer(),
realm :: binary(),
identity :: binary(),
geo_identity :: map(), %% geo metadata for CONNECT (city, country, lat, lng)
site :: map() | undefined, %% site metadata for CONNECT (site_id, name, city, lat, lng, site_type)
client_type :: binary(), %% <<"node">> or <<"relay">> (for CONNECT handshake)
tls_verify :: verify_peer | none, %% TLS verification mode (none for relay peering)
previous_host :: string() | undefined, %% host before failover (for reroute detection)
conn :: reference() | undefined,
stream :: reference() | undefined,
status :: connecting | connected | disconnected,
recv_buffer :: binary(),
reconnect_attempt :: non_neg_integer(), %% for exponential backoff
%% Application state (survives reconnects)
subscriptions :: #{reference() => {binary(), fun()}}, %% ref => {topic, callback}
procedures :: #{binary() => fun()}, %% procedure => handler
pending_calls :: #{binary() => {pid(), reference()}}, %% call_id => {from, timer_ref}
ping_sent_at :: integer() | undefined, %% erlang:monotonic_time(millisecond) when PING sent
last_rtt_ms :: non_neg_integer() | undefined %% most recent RTT measurement
}).
%%====================================================================
%% API
%%====================================================================
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
-spec stop(pid()) -> ok.
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).
%%====================================================================
%% 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,
%% Randomize initial relay to distribute load across relays
Index = rand:uniform(length(Relays)) - 1,
Url = lists:nth(Index + 1, Relays),
{Host, Port} = parse_url(Url),
Realm = maps:get(realm, Opts, <<"io.macula">>),
Identity = maps:get(identity, Opts, <<"anonymous">>),
GeoIdentity = maps:get(geo_identity, Opts, #{}),
Site = maps:get(site, Opts, undefined),
ClientType = maps:get(type, Opts, <<"node">>),
TlsVerify = maps:get(tls_verify, Opts, verify_peer),
State = #state{
relays = Relays,
relay_index = Index,
url = Url,
host = Host,
port = Port,
realm = Realm,
identity = Identity,
geo_identity = GeoIdentity,
site = Site,
client_type = ClientType,
tls_verify = TlsVerify,
conn = undefined,
stream = undefined,
status = connecting,
recv_buffer = <<>>,
reconnect_attempt = 0,
subscriptions = #{},
procedures = #{},
pending_calls = #{}
},
self() ! connect,
{ok, State}.
%%====================================================================
%% Subscribe / Unsubscribe
%%====================================================================
handle_call({subscribe, Topic, Callback}, _From, State) ->
Ref = make_ref(),
Subs = maps:put(Ref, {Topic, Callback}, State#state.subscriptions),
%% Send SUBSCRIBE if connected
maybe_send(subscribe, #{<<"topics">> => [Topic], <<"qos">> => 0}, State),
{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),
maybe_send(unsubscribe, #{<<"topics">> => [Topic]}, State),
{reply, ok, State#state{subscriptions = Subs}}
end;
%%====================================================================
%% Advertise / Unadvertise
%%====================================================================
handle_call({advertise, Procedure, Handler}, _From, State) ->
Procs = maps:put(Procedure, Handler, State#state.procedures),
maybe_send(register_procedure, #{<<"procedure">> => Procedure}, State),
{reply, {ok, make_ref()}, State#state{procedures = Procs}};
handle_call({unadvertise, Procedure}, _From, State) ->
Procs = maps:remove(Procedure, State#state.procedures),
%% No protocol message for unadvertise yet — procedure removed on disconnect
{reply, ok, State#state{procedures = Procs}};
%%====================================================================
%% RPC Call
%%====================================================================
handle_call({rpc_call, Procedure, Args, Timeout}, From, State) ->
CallId = base64:encode(crypto:strong_rand_bytes(12)),
TimerRef = erlang:send_after(Timeout, self(), {call_timeout, CallId}),
PendingCalls = maps:put(CallId, {From, TimerRef}, State#state.pending_calls),
maybe_send(call, #{
<<"call_id">> => CallId,
<<"procedure">> => Procedure,
<<"args">> => Args
}, State),
{noreply, State#state{pending_calls = PendingCalls}};
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
%%====================================================================
%% Publish (async)
%%====================================================================
handle_cast({publish, Topic, Payload}, State) ->
BinPayload = to_binary_payload(Payload),
maybe_send(publish, #{
<<"topic">> => Topic,
<<"payload">> => BinPayload,
<<"qos">> => 0,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
}, State),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%%====================================================================
%% Connection lifecycle
%%====================================================================
handle_info(connect, State) ->
QuicOpts = [{alpn, ["macula"]} | build_tls_opts(State#state.tls_verify)],
case macula_quic:connect(State#state.host, State#state.port, QuicOpts, 10000) of
{ok, Conn} ->
case macula_quic:open_stream(Conn) of
{ok, Stream} ->
quicer:setopt(Stream, active, true),
?LOG_INFO("[relay_client] Connected to ~s", [State#state.url]),
State2 = State#state{conn = Conn, stream = Stream, status = connected,
recv_buffer = <<>>, reconnect_attempt = 0},
send_connect(State2),
replay_state(State2),
schedule_ping(),
{noreply, State2};
{error, StreamErr} ->
?LOG_WARNING("[relay_client] Stream open failed on ~s: ~p", [State#state.url, StreamErr]),
catch macula_quic:close(Conn),
{noreply, schedule_failover(State)}
end;
{error, Reason} ->
?LOG_WARNING("[relay_client] Connect to ~s failed: ~p", [State#state.url, Reason]),
{noreply, schedule_failover(State)};
{error, Type, Detail} ->
?LOG_WARNING("[relay_client] Connect to ~s failed: ~p ~p", [State#state.url, Type, Detail]),
{noreply, schedule_failover(State)}
end;
handle_info({call_timeout, CallId}, State) ->
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined -> {noreply, State};
{From, _TimerRef} ->
gen_server:reply(From, {error, timeout}),
{noreply, State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)}}
end;
%%====================================================================
%% QUIC data received
%%====================================================================
handle_info({quic, Data, _Stream, _Flags}, State) when is_binary(Data) ->
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
{NewBuffer, State2} = process_buffer(Buffer, State),
{noreply, State2#state{recv_buffer = NewBuffer}};
%%====================================================================
%% QUIC lifecycle — reconnect on any closure
%%====================================================================
handle_info({quic, peer_send_shutdown, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, peer_send_aborted, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, send_shutdown_complete, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, shutdown, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({quic, closed, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({quic, transport_shutdown, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({error, transport_down, _Detail}, State) ->
handle_disconnect(State);
handle_info({quic, streams_available, _Conn, _Info}, State) ->
{noreply, State};
handle_info({quic, peer_needs_streams, _Conn, _Info}, State) ->
{noreply, State};
handle_info(send_ping, #state{status = connected} = State) ->
send_protocol_ping(State),
schedule_ping(),
{noreply, State#state{ping_sent_at = erlang:monotonic_time(millisecond)}};
handle_info(send_ping, State) ->
%% Not connected — skip, next ping will fire after reconnect
{noreply, State};
handle_info(Info, State) ->
?LOG_DEBUG("[relay_client] Unhandled: ~p", [Info]),
{noreply, State}.
terminate(_Reason, #state{conn = Conn, stream = Stream}) ->
catch macula_quic:close(Stream),
catch macula_quic:close(Conn),
ok.
%%====================================================================
%% Mesh ping — RTT measurement + publish
%%====================================================================
schedule_ping() ->
erlang:send_after(?PING_INTERVAL_MS, self(), send_ping).
send_protocol_ping(State) ->
Ts = erlang:system_time(millisecond),
maybe_send(ping, #{timestamp => Ts}, State).
publish_mesh_ping(RttMs, State) ->
RelayHost = list_to_binary(State#state.host),
NodeName = State#state.identity,
Site = State#state.site,
Lat = site_field(lat, Site),
Lng = site_field(lng, Site),
Payload = json:encode(#{
relay => RelayHost,
node => NodeName,
rtt_ms => RttMs,
lat => Lat,
lng => Lng,
ts => erlang:system_time(second)
}),
maybe_send(publish, #{
<<"topic">> => ?PING_TOPIC,
<<"payload">> => iolist_to_binary(Payload),
<<"qos">> => 0,
<<"retain">> => false,
<<"message_id">> => base64:encode(crypto:strong_rand_bytes(8))
}, State).
notify_discovery(RttMs, State) ->
RelayHost = list_to_binary(State#state.host),
try macula_relay_discovery ! {ping_rtt, RelayHost, RttMs}
catch _:_ -> ok
end.
site_field(Key, Site) when is_map(Site) -> maps:get(Key, Site, 0.0);
site_field(_, _) -> 0.0.
%%====================================================================
%% Internal: message processing
%%====================================================================
process_buffer(Buffer, State) when byte_size(Buffer) < 8 ->
{Buffer, State};
process_buffer(<<_Version:8, _TypeId:8, _Flags:8, _Reserved:8,
PayloadLen:32/big-unsigned, Rest/binary>> = Buffer, State) ->
case byte_size(Rest) >= PayloadLen of
true ->
MsgBytes = binary:part(Buffer, 0, 8 + PayloadLen),
Remaining = binary:part(Buffer, 8 + PayloadLen, byte_size(Buffer) - 8 - PayloadLen),
State2 = handle_message(macula_protocol_decoder:decode(MsgBytes), State),
process_buffer(Remaining, State2);
false ->
{Buffer, State}
end.
%% Incoming PUBLISH from relay
handle_message({ok, {publish, Msg}}, State) ->
Topic = maps:get(<<"topic">>, Msg, <<>>),
Payload = maps:get(<<"payload">>, Msg, <<>>),
invoke_matching_callbacks(Topic, Payload, State#state.subscriptions),
State;
%% Incoming CALL from relay (we're the provider)
handle_message({ok, {call, Msg}}, State) ->
Procedure = maps:get(<<"procedure">>, Msg, <<>>),
CallId = maps:get(<<"call_id">>, Msg, <<>>),
Args = maps:get(<<"args">>, Msg, #{}),
case maps:get(Procedure, State#state.procedures, undefined) of
undefined ->
maybe_send(reply, #{<<"call_id">> => CallId,
<<"error">> => #{<<"code">> => <<"not_found">>}}, State);
Handler ->
%% Execute handler in separate process to avoid blocking
Stream = State#state.stream,
spawn(fun() ->
Result = try Handler(Args)
catch _:Err -> {error, Err}
end,
ReplyMsg = case Result of
{ok, R} -> #{<<"call_id">> => CallId, <<"result">> => R};
{error, R} -> #{<<"call_id">> => CallId,
<<"error">> => #{<<"code">> => <<"handler_error">>,
<<"message">> => iolist_to_binary(
io_lib:format("~p", [R]))}}
end,
Binary = macula_protocol_encoder:encode(reply, ReplyMsg),
macula_quic:async_send(Stream, Binary)
end)
end,
State;
%% Incoming REPLY from relay (we're the caller)
handle_message({ok, {reply, Msg}}, State) ->
CallId = maps:get(<<"call_id">>, Msg, <<>>),
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined ->
?LOG_DEBUG("[relay_client] Reply for unknown call_id"),
State;
{From, TimerRef} ->
erlang:cancel_timer(TimerRef),
Result = case maps:get(<<"result">>, Msg, undefined) of
undefined ->
Error = maps:get(<<"error">>, Msg, #{}),
{error, Error};
R ->
{ok, R}
end,
gen_server:reply(From, Result),
State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)}
end;
%% PONG — measure RTT and publish mesh ping
handle_message({ok, {pong, _Msg}}, #state{ping_sent_at = undefined} = State) ->
State;
handle_message({ok, {pong, _Msg}}, State) ->
RttMs = erlang:monotonic_time(millisecond) - State#state.ping_sent_at,
publish_mesh_ping(RttMs, State),
notify_discovery(RttMs, State),
State#state{ping_sent_at = undefined, last_rtt_ms = RttMs};
handle_message({ok, {Type, _Msg}}, State) ->
?LOG_DEBUG("[relay_client] Ignoring message type: ~p", [Type]),
State;
handle_message({error, Reason}, State) ->
?LOG_WARNING("[relay_client] Decode error: ~p", [Reason]),
State.
%%====================================================================
%% Internal helpers
%%====================================================================
send_connect(State) ->
NodeId = crypto:strong_rand_bytes(32),
%% Use identity as node_name when set (multi-tenant stubs).
%% Falls back to Erlang node name, then random hex.
NodeName = case State#state.identity of
Id when is_binary(Id), byte_size(Id) > 0 -> Id;
_ ->
case net_kernel:nodename() of
nonode@nohost -> binary:encode_hex(crypto:strong_rand_bytes(8));
Name -> atom_to_binary(Name)
end
end,
%% target_relay = the hostname the client connected to (for multi-tenant relay routing)
TargetRelay = list_to_binary(State#state.host),
Base = #{
version => <<"1.0">>,
node_id => NodeId,
node_name => NodeName,
realm_id => State#state.realm,
capabilities => [<<"pubsub">>, <<"rpc">>],
endpoint => State#state.url,
type => State#state.client_type,
target_relay => TargetRelay
},
Msg0 = case State#state.geo_identity of
GeoId when is_map(GeoId), map_size(GeoId) > 0 -> Base#{identity => GeoId};
_ ->
case macula_connection:build_node_identity(#{}) of
EnvId when map_size(EnvId) > 0 -> Base#{identity => EnvId};
_ -> Base
end
end,
%% Include site metadata if configured
Msg1 = case State#state.site of
SiteMap when is_map(SiteMap), map_size(SiteMap) > 0 -> Msg0#{site => SiteMap};
_ -> Msg0
end,
%% Include previous_relay if this is a failover reconnect (different host)
Msg = case State#state.previous_host of
PrevHost when is_list(PrevHost), PrevHost =/= State#state.host ->
Msg1#{previous_relay => list_to_binary(PrevHost)};
_ -> Msg1
end,
maybe_send(connect, Msg, State).
%% Replay all subscriptions and procedure registrations after reconnect
replay_state(#state{subscriptions = Subs, procedures = Procs} = State) ->
%% Batch all topics in one SUBSCRIBE
Topics = lists:usort([T || {_Ref, {T, _Cb}} <- maps:to_list(Subs)]),
case Topics of
[] -> ok;
_ ->
?LOG_INFO("[relay_client] Replaying ~p subscription(s)", [length(Topics)]),
maybe_send(subscribe, #{<<"topics">> => Topics, <<"qos">> => 0}, State)
end,
%% Register all procedures
maps:foreach(fun(Proc, _Handler) ->
?LOG_INFO("[relay_client] Replaying procedure: ~s", [Proc]),
maybe_send(register_procedure, #{<<"procedure">> => Proc}, State)
end, Procs).
maybe_send(_Type, _Msg, #state{stream = undefined}) -> ok;
maybe_send(Type, Msg, #state{stream = Stream}) ->
Binary = macula_protocol_encoder:encode(Type, Msg),
macula_quic:async_send(Stream, Binary).
%% Ensure payload is a flat binary before passing to msgpack.
%% Callers may pass: map, iolist (from json:encode), or binary.
to_binary_payload(Payload) when is_binary(Payload) -> Payload;
to_binary_payload(Payload) when is_map(Payload) -> iolist_to_binary(json:encode(Payload));
to_binary_payload(Payload) when is_list(Payload) -> iolist_to_binary(Payload);
to_binary_payload(Payload) -> iolist_to_binary(json:encode(Payload)).
handle_disconnect(#state{status = disconnected} = State) ->
{noreply, State};
handle_disconnect(State) ->
?LOG_WARNING("[relay_client] Disconnected from ~s", [State#state.url]),
catch macula_quic:close(State#state.stream),
catch macula_quic:close(State#state.conn),
State2 = State#state{conn = undefined, stream = undefined, status = disconnected,
recv_buffer = <<>>},
{noreply, schedule_failover(State2)}.
%% Pick next relay and schedule reconnect with exponential backoff + jitter.
%% Uses geographic discovery if available, falls back to round-robin.
schedule_failover(#state{reconnect_attempt = Attempt} = State) ->
{NextUrl, NextState} = select_failover_relay(State),
{NextHost, NextPort} = parse_url(NextUrl),
BackoffMs = backoff_ms(Attempt),
erlang:send_after(BackoffMs, self(), connect),
?LOG_INFO("[relay_client] Failover to ~s in ~bms (attempt ~b)",
[NextUrl, BackoffMs, Attempt + 1]),
NextState#state{url = NextUrl, host = NextHost, port = NextPort,
status = disconnected, reconnect_attempt = Attempt + 1,
previous_host = State#state.host}.
%% Try geographic discovery first, fall back to round-robin.
select_failover_relay(#state{host = CurrentHost} = State) ->
CurrentHostname = list_to_binary(CurrentHost),
case try_discovery_failover(CurrentHostname) of
{ok, Url} ->
{Url, State};
{error, _} ->
round_robin_failover(State)
end.
try_discovery_failover(CurrentHostname) ->
try macula_relay_discovery:nearest_except(CurrentHostname)
catch _:_ -> {error, discovery_unavailable}
end.
round_robin_failover(#state{relays = Relays, relay_index = CurrentIdx} = State) ->
NumRelays = length(Relays),
NextIdx = case NumRelays > 1 of
true -> (CurrentIdx + 1) rem NumRelays;
false -> CurrentIdx
end,
NextUrl = lists:nth(NextIdx + 1, Relays),
{NextUrl, State#state{relay_index = NextIdx}}.
backoff_ms(Attempt) ->
Base = min(?RECONNECT_BASE_MS * (1 bsl min(Attempt, 10)), ?RECONNECT_MAX_MS),
Jitter = round(Base * ?RECONNECT_JITTER),
Base + rand:uniform(2 * Jitter + 1) - Jitter - 1.
parse_url(Url) when is_binary(Url) ->
parse_url(binary_to_list(Url));
parse_url("https://" ++ Rest) ->
parse_host_port(Rest);
parse_url(Rest) ->
parse_host_port(Rest).
parse_host_port(HostPort) ->
case string:split(HostPort, ":", trailing) of
[Host, PortStr] -> {Host, list_to_integer(PortStr)};
[Host] -> {Host, 4433}
end.
invoke_matching_callbacks(Topic, Payload, Subscriptions) ->
DecodedPayload = safe_decode_json(Payload),
maps:foreach(fun(_Ref, {SubTopic, Callback}) ->
case topic_matches(SubTopic, Topic) of
true ->
spawn(fun() ->
try Callback(#{topic => Topic, payload => DecodedPayload})
catch C:E -> ?LOG_WARNING("[relay_client] Callback error: ~p:~p", [C, E])
end
end);
false -> ok
end
end, Subscriptions).
topic_matches(Pattern, Topic) when Pattern =:= Topic -> true;
topic_matches(<<"**">>, _Topic) -> true;
topic_matches(Pattern, Topic) ->
PatternParts = binary:split(Pattern, <<".">>, [global]),
TopicParts = binary:split(Topic, <<".">>, [global]),
parts_match(PatternParts, TopicParts).
parts_match([], []) -> true;
parts_match([<<"**">>], _) -> true;
parts_match([<<"*">> | PR], [_ | TR]) -> parts_match(PR, TR);
parts_match([P | PR], [T | TR]) when P =:= T -> parts_match(PR, TR);
parts_match(_, _) -> false.
safe_decode_json(Payload) when is_binary(Payload) ->
try json:decode(Payload) catch _:_ -> Payload end;
safe_decode_json(Payload) -> Payload.
%% Build TLS options for QUIC client connection.
%% verify_peer: use global TLS mode (production or development).
%% none: force no-verify opts regardless of global mode.
build_tls_opts(none) ->
[{verify, none}];
build_tls_opts(_) ->
macula_tls:quic_client_opts().