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
Current section
Files
src/macula_relay.erl
%%%-------------------------------------------------------------------
%%% @doc Macula Relay Server — hub-spoke message routing via pg + gproc.
%%%
%%% Accepts QUIC connections from nodes. Each connection gets a handler
%%% process. Pub/sub uses OTP pg groups. RPC uses gproc registry.
%%% Process death = automatic cleanup (no TTLs, no manual eviction).
%%%
%%% Uses async accept pattern (like macula_gateway_quic_server) —
%%% quicer:async_accept delivers {quic, new_conn, Conn, Info} messages.
%%%
%%% Start: `macula_relay:start_link(#{port => 4433}).'
%%% @end
%%%-------------------------------------------------------------------
-module(macula_relay).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([start_link/1, start_link/2]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
listener :: reference() | undefined,
port :: integer(),
handlers :: [pid()],
%% stream_ref => {handler_pid, [buffered_data]}
%% Captures data that arrives before ownership transfer completes
stream_handlers :: #{reference() => pid()}
}).
%%====================================================================
%% API
%%====================================================================
%% @doc Start the default (singleton) relay listener.
start_link(Opts) ->
gen_server:start_link({local, ?MODULE}, ?MODULE, Opts, []).
%% @doc Start a named relay listener (for per-identity binding).
%% Name is typically the identity hostname as an atom.
start_link(Name, Opts) ->
gen_server:start_link({local, Name}, ?MODULE, Opts, []).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init(Opts) ->
Port = maps:get(port, Opts, 4433),
BindAddr = maps:get(bind_addr, Opts, undefined),
%% Ensure pg scope exists
case pg:start(pg) of
{ok, _} -> ok;
{error, {already_started, _}} -> ok
end,
%% Get TLS certs — check env vars first (production), then macula_tls (dev)
{CertPath, KeyPath} = get_tls_paths(),
?LOG_INFO("[relay] TLS cert: ~s, key: ~s", [CertPath, KeyPath]),
ListenOpts = [
{cert, CertPath},
{key, KeyPath},
{alpn, ["macula"]},
{peer_unidi_stream_count, 3},
{peer_bidi_stream_count, 100},
{idle_timeout_ms, 120000},
{keep_alive_interval_ms, 30000}
],
ListenTarget = case BindAddr of
undefined -> Port;
Addr -> {binary_to_list(Addr), Port}
end,
case macula_quic:listen(ListenTarget, ListenOpts) of
{ok, Listener} ->
?LOG_INFO("[relay] Listening on ~p port ~p", [BindAddr, Port]),
%% Register for async connection events
register_accept(Listener),
{ok, #state{listener = Listener, port = Port, handlers = [],
stream_handlers = #{}}};
{error, Reason} ->
{stop, {listen_failed, Reason}}
end.
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% Async accept: quicer delivers new connections as messages
handle_info({quic, new_conn, Conn, ConnInfo}, #state{listener = Listener} = State) ->
?LOG_INFO("[relay] New connection: ~p", [ConnInfo]),
complete_handshake(Conn),
%% MUST re-register for next connection
register_accept(Listener),
{noreply, State};
%% Async stream: quicer delivers new streams as messages
handle_info({quic, new_stream, Stream, StreamProps}, State) ->
?LOG_INFO("[relay] New stream: ~p", [StreamProps]),
%% Ensure stream is passive until handler takes over — prevents
%% data arriving on the relay process before ownership transfer.
catch quicer:setopt(Stream, active, false),
Conn = case StreamProps of
#{conn := C} -> C;
_ -> get(pending_conn)
end,
case Conn of
undefined ->
?LOG_WARNING("[relay] Stream arrived with no connection context"),
{noreply, State};
_ ->
{ok, Pid} = macula_relay_handler:start_link(Conn, Stream),
case quicer:controlling_process(Stream, Pid) of
ok ->
%% Connection ownership may fail if another handler already owns it
%% (race with concurrent streams on same connection). That's OK —
%% the handler still works with stream ownership.
_ = quicer:controlling_process(Conn, Pid),
Pid ! ownership_transferred,
?LOG_INFO("[relay] Handler ~p started, ownership transferred", [Pid]),
SH = maps:put(Stream, Pid, State#state.stream_handlers),
{noreply, State#state{handlers = [Pid | State#state.handlers],
stream_handlers = SH}};
{error, Reason} ->
?LOG_WARNING("[relay] Failed to transfer stream ownership: ~p", [Reason]),
{noreply, State}
end
end;
%% QUIC data arriving on relay process before handler takes over.
%% Forward to the handler that owns this stream.
handle_info({quic, Data, Stream, Flags}, State) when is_binary(Data) ->
case maps:get(Stream, State#state.stream_handlers, undefined) of
undefined ->
?LOG_WARNING("[relay] Data for unknown stream (~p bytes)", [byte_size(Data)]);
Pid ->
Pid ! {quic, Data, Stream, Flags}
end,
{noreply, State};
%% Handler process died — remove from handlers list and stream map
handle_info({'EXIT', Pid, Reason}, State) ->
case lists:member(Pid, State#state.handlers) of
true ->
?LOG_INFO("[relay] Handler ~p exited: ~p", [Pid, Reason]),
SH = maps:filter(fun(_S, P) -> P =/= Pid end, State#state.stream_handlers),
{noreply, State#state{handlers = lists:delete(Pid, State#state.handlers),
stream_handlers = SH}};
false ->
{noreply, State}
end;
handle_info(Info, State) ->
?LOG_DEBUG("[relay] Unhandled info: ~p", [Info]),
{noreply, State}.
terminate(_Reason, #state{listener = Listener}) ->
catch macula_quic:close(Listener),
ok.
%%====================================================================
%% Internal
%%====================================================================
%% Register for async accept — quicer will send {quic, new_conn, ...}
register_accept(Listener) ->
case quicer:async_accept(Listener, #{}) of
{ok, _} ->
?LOG_INFO("[relay] Ready for connections"),
ok;
{error, Reason} ->
?LOG_WARNING("[relay] async_accept failed: ~p", [Reason]),
ok
end.
%% Complete TLS handshake, then register for async stream accept
complete_handshake(Conn) ->
case quicer:handshake(Conn) of
ok ->
accept_streams(Conn);
{ok, _} ->
accept_streams(Conn);
{error, Reason} ->
?LOG_ERROR("[relay] Handshake failed: ~p", [Reason]),
catch quicer:close_connection(Conn)
end.
accept_streams(Conn) ->
?LOG_INFO("[relay] Handshake complete, accepting streams"),
put(pending_conn, Conn),
case quicer:async_accept_stream(Conn, #{}) of
{ok, _} ->
ok;
{error, Reason} ->
?LOG_WARNING("[relay] async_accept_stream failed: ~p", [Reason]),
ok
end.
%% @private Get TLS cert/key paths. Env vars take precedence (production),
%% then app config, then auto-generate for dev.
get_tls_paths() ->
CertPath = case os:getenv("MACULA_TLS_CERTFILE") of
false -> get_tls_path_from_config(cert);
C -> C
end,
KeyPath = case os:getenv("MACULA_TLS_KEYFILE") of
false -> get_tls_path_from_config(key);
K -> K
end,
{CertPath, KeyPath}.
get_tls_path_from_config(cert) ->
{CertPath, _} = macula_tls:get_cert_paths(),
{ok, CertPath2, _, _} = macula_tls:ensure_cert_exists(CertPath, element(2, macula_tls:get_cert_paths())),
CertPath2;
get_tls_path_from_config(key) ->
{_, KeyPath} = macula_tls:get_cert_paths(),
{ok, _, KeyPath2, _} = macula_tls:ensure_cert_exists(element(1, macula_tls:get_cert_paths()), KeyPath),
KeyPath2.