Packages
macula
3.3.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_dist_system/macula_dist_relay.erl
%%%-------------------------------------------------------------------
%%% @doc Relay-routed Erlang distribution.
%%%
%%% When MACULA_DIST_MODE=relay, Erlang distribution connections
%%% are tunneled through the Macula relay mesh instead of direct QUIC.
%%% This enables distribution across firewalls and NATs — nodes only
%%% need outbound connectivity to a relay.
%%%
%%% The tunnel works by:
%%% 1. Node A requests a tunnel via mesh RPC (_dist.tunnel.{node})
%%% 2. Node B creates a gen_tcp loopback pair bridged to relay pub/sub
%%% 3. Node A creates its own loopback pair bridged to the tunnel
%%% 4. dist_util handshake flows through the relay tunnel
%%% 5. Post-handshake: tick keepalive + distribution traffic via bridge
%%%
%%% Bridge processes are supervised by `macula_dist_bridge_sup'
%%% (simple_one_for_one under `macula_dist_system').
%%%
%%% Tunnel bytes are encrypted with AES-256-GCM derived from the
%%% Erlang distribution cookie. The relay cannot read ETF content.
%%%
%%% == Recommended net_ticktime ==
%%%
%%% Relay adds WAN latency to every tick. Increase net_ticktime
%%% to avoid false disconnects:
%%% -kernel net_ticktime 120
%%% The default 60s may cause spurious node-DOWN on high-latency relays.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_dist_relay).
-include_lib("kernel/include/logger.hrl").
-export([connect/3]).
-export([is_relay_mode/0, get_mesh_client/0]).
-export([register_mesh_client/1]).
-export([advertise_dist_accept/0]).
-export([extract_payload/1]).
-export([get_tunnel_metrics/0, get_tunnel_metrics/1]).
%% Must be shorter than OTP's SetupTime. Default SetupTime is 7000ms
%% but we set it to 15000ms via dist_setup_timeout (see macula_dist.erl).
%% If the tunnel RPC takes longer than DIST_TIMEOUT, the caller gets
%% {error, timeout} with diagnostics. If it takes longer than
%% SetupTime, dist_util kills do_setup with zero diagnostics (just pang).
-define(DIST_TIMEOUT, 10000).
-define(CONTROLLER_TIMEOUT, 30000).
-define(METRIC_BYTES_OUT, 1).
-define(METRIC_BYTES_IN, 2).
-define(METRIC_MSGS_OUT, 3).
-define(METRIC_MSGS_IN, 4).
%%%===================================================================
%%% Public API
%%%===================================================================
%% @doc Register a mesh relay client for distribution tunneling.
-spec register_mesh_client(pid()) -> ok.
register_mesh_client(Pid) when is_pid(Pid) ->
persistent_term:put(macula_dist_mesh_client, Pid),
?LOG_INFO("[dist_relay] Mesh client registered: ~p", [Pid]),
ok.
%% @doc Check if relay distribution mode is enabled.
-spec is_relay_mode() -> boolean().
is_relay_mode() ->
os:getenv("MACULA_DIST_MODE") =:= "relay".
%% @doc Connect to a remote node via relay mesh.
-spec connect(string(), string(), integer()) -> {ok, port(), port()} | {error, term()}.
connect(NodeStr, _Host, _Port) ->
?LOG_INFO("[dist_relay] Connecting to ~s via mesh", [NodeStr]),
case get_mesh_client() of
undefined ->
{error, no_mesh_connection};
MeshClient ->
request_tunnel(MeshClient, NodeStr)
end.
%% @doc Advertise this node as accepting distribution connections via relay.
-spec advertise_dist_accept() -> ok.
advertise_dist_accept() ->
ensure_bridge_sup(),
case get_mesh_client() of
undefined ->
?LOG_WARNING("[dist_relay] Cannot advertise — no mesh client registered"),
ok;
Client ->
NodeName = atom_to_binary(node()),
Procedure = <<"_dist.tunnel.", NodeName/binary>>,
Handler = fun(Args) -> handle_tunnel_request(Args) end,
macula_mesh_client:advertise(Client, Procedure, Handler),
?LOG_INFO("[dist_relay] Advertised distribution accept: ~s", [Procedure]),
ok
end.
%% Bridge supervisor is started by macula_root (application supervisor),
%% so it survives shell crashes and user code failures. This is a no-op
%% kept for backwards compatibility.
ensure_bridge_sup() ->
case whereis(macula_dist_bridge_sup) of
undefined ->
?LOG_WARNING("[dist_relay] Bridge sup not running — macula app may not be started");
_Pid ->
ok
end.
%% @doc Get metrics for all active tunnels.
-spec get_tunnel_metrics() -> [{binary(), map()}].
get_tunnel_metrics() ->
case persistent_term:get(macula_dist_tunnels, undefined) of
undefined -> [];
Tunnels ->
maps:fold(fun(TunnelId, Ref, Acc) ->
[{TunnelId, read_metrics(Ref)} | Acc]
end, [], Tunnels)
end.
%% @doc Get metrics for a specific tunnel.
-spec get_tunnel_metrics(binary()) -> map() | undefined.
get_tunnel_metrics(TunnelId) ->
case persistent_term:get(macula_dist_tunnels, undefined) of
undefined -> undefined;
Tunnels ->
case maps:get(TunnelId, Tunnels, undefined) of
undefined -> undefined;
Ref -> read_metrics(Ref)
end
end.
%%%===================================================================
%%% Internal — Mesh Client Lookup
%%%===================================================================
-spec get_mesh_client() -> pid() | undefined.
get_mesh_client() ->
case persistent_term:get(macula_dist_mesh_client, undefined) of
undefined -> undefined;
Pid when is_pid(Pid) ->
case is_process_alive(Pid) of
true -> Pid;
false -> undefined
end
end.
%%%===================================================================
%%% Internal — Tunnel Negotiation (connecting side)
%%%===================================================================
request_tunnel(MeshClient, NodeStr) ->
Procedure = <<"_dist.tunnel.", (list_to_binary(NodeStr))/binary>>,
Args = #{<<"from_node">> => atom_to_binary(node()),
<<"target_node">> => list_to_binary(NodeStr)},
?LOG_INFO("[dist_relay] RPC ~s via ~p", [Procedure, MeshClient]),
RawResult = tunnel_rpc(MeshClient, Procedure, Args),
%% Normalize: mesh_client may return {ok, R, TraceMap} when relay
%% includes _trace in the call round-trip. Strip the trace wrapper
%% so pattern matching below always sees a 2-tuple.
Result = normalize_rpc_result(RawResult),
case Result of
{ok, #{<<"tunnel_id">> := TunnelId}} ->
?LOG_INFO("[dist_relay] Tunnel established: ~s", [TunnelId]),
create_dist_socket(MeshClient, TunnelId);
{ok, #{<<"error">> := ErrorInfo}} ->
?LOG_WARNING("[dist_relay] Tunnel error: ~p", [ErrorInfo]),
{error, {tunnel_error, ErrorInfo}};
{error, Reason} ->
?LOG_WARNING("[dist_relay] Tunnel request failed: ~p", [Reason]),
{error, {tunnel_failed, Reason}};
Other ->
?LOG_WARNING("[dist_relay] Unexpected RPC result: ~p", [Other]),
{error, {unexpected_result, Other}}
end.
normalize_rpc_result({ok, R, _Trace}) -> {ok, R};
normalize_rpc_result(Other) -> Other.
%% Use call_any for multi_relay (tries each connected relay).
%% Fall back to direct call for single relay_client.
tunnel_rpc(MeshClient, Procedure, Args) ->
case is_multi_relay(MeshClient) of
true ->
macula_multi_relay:call_any(MeshClient, Procedure, Args, ?DIST_TIMEOUT);
false ->
macula_mesh_client:call(MeshClient, Procedure, Args, ?DIST_TIMEOUT)
end.
%% Check if the PID is a macula_multi_relay gen_server by inspecting
%% the registered name or the initial call in process_info.
is_multi_relay(Pid) ->
case erlang:process_info(Pid, dictionary) of
{dictionary, Dict} ->
case proplists:get_value('$initial_call', Dict) of
{macula_multi_relay, init, 1} -> true;
_ -> false
end;
_ ->
false
end.
%%%===================================================================
%%% Internal — Tunnel Negotiation (accepting side)
%%%===================================================================
handle_tunnel_request(Args) ->
FromNode = maps:get(<<"from_node">>, Args, <<>>),
TunnelId = base64:encode(crypto:strong_rand_bytes(16)),
?LOG_INFO("[dist_relay] Tunnel request from ~s, id: ~s", [FromNode, TunnelId]),
SendTopic = <<"_dist.data.", TunnelId/binary, ".in">>,
RecvTopic = <<"_dist.data.", TunnelId/binary, ".out">>,
case get_mesh_client() of
undefined ->
{error, <<"no_mesh_client">>};
Client ->
spawn_accept_bridge(Client, TunnelId, SendTopic, RecvTopic),
{ok, #{<<"tunnel_id">> => TunnelId,
<<"send_topic">> => SendTopic,
<<"recv_topic">> => RecvTopic}}
end.
%% The accept side has a setup phase (kernel negotiation) before the
%% supervised bridge can start. This setup process is short-lived —
%% it creates the loopback pair, negotiates with net_kernel, then
%% hands off to the supervised bridge.
spawn_accept_bridge(Client, TunnelId, SendTopic, RecvTopic) ->
{Pid, _Ref} = spawn_monitor(fun() ->
dist_accept_setup(Client, TunnelId, SendTopic, RecvTopic)
end),
?LOG_INFO("[dist_relay] Accept setup ~p for ~s", [Pid, TunnelId]),
Pid.
%%%===================================================================
%%% Internal — Accept-side Setup (short-lived, then hands to bridge)
%%%===================================================================
dist_accept_setup(Client, TunnelId, SendTopic, RecvTopic) ->
Self = self(),
Key = tunnel_key(),
%% Subscribe for buffering before kernel knows about us
{ok, BufferSubRef} = macula_mesh_client:subscribe(Client, RecvTopic,
fun(Msg) -> Self ! {dist_data_buffered, extract_payload(Msg)} end),
{DistSock, BridgeSock} = create_loopback_pair(),
case whereis(net_kernel) of
undefined ->
?LOG_WARNING("[dist_bridge] net_kernel not found"),
macula_mesh_client:unsubscribe(Client, BufferSubRef);
KernelPid ->
negotiate_with_kernel(Self, KernelPid, Client, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId, Key)
end.
negotiate_with_kernel(Self, KernelPid, Client, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId, Key) ->
KernelPid ! {accept, Self, DistSock, inet, macula_dist},
receive
{KernelPid, controller, DistCtrl} ->
DistCtrl ! {Self, controller, ok},
%% Transfer DistSock to the dist controller so it survives
%% when this setup process exits.
gen_tcp:controlling_process(DistSock, DistCtrl),
macula_mesh_client:unsubscribe(Client, BufferSubRef),
flush_buffered_to_socket(BridgeSock, Key),
start_supervised_bridge(Client, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key);
{KernelPid, unsupported_protocol} ->
?LOG_WARNING("[dist_bridge] Unsupported protocol"),
macula_mesh_client:unsubscribe(Client, BufferSubRef)
after ?CONTROLLER_TIMEOUT ->
?LOG_WARNING("[dist_bridge] Controller timeout"),
macula_mesh_client:unsubscribe(Client, BufferSubRef)
end.
%%%===================================================================
%%% Internal — Connect-side Bridge Creation
%%%===================================================================
create_dist_socket(MeshClient, TunnelId) ->
SendTopic = <<"_dist.data.", TunnelId/binary, ".out">>,
RecvTopic = <<"_dist.data.", TunnelId/binary, ".in">>,
{DistSock, BridgeSock} = create_loopback_pair(),
Key = tunnel_key(),
start_supervised_bridge(MeshClient, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key),
{ok, DistSock, DistSock}.
%%%===================================================================
%%% Internal — Supervised Bridge Startup
%%%===================================================================
start_supervised_bridge(Client, BridgeSock, SendTopic, RecvTopic, TunnelId, Key) ->
Metrics = init_metrics(TunnelId),
BridgeArgs = #{
client => Client,
bridge_sock => BridgeSock,
send_topic => SendTopic,
recv_topic => RecvTopic,
tunnel_id => TunnelId,
key => Key,
metrics => Metrics
},
case macula_dist_bridge_sup:start_bridge(BridgeArgs) of
{ok, Pid} ->
gen_tcp:controlling_process(BridgeSock, Pid),
Pid ! socket_ready,
?LOG_INFO("[dist_relay] Supervised bridge ~p for ~s", [Pid, TunnelId]),
{ok, Pid};
{error, Reason} ->
?LOG_ERROR("[dist_relay] Failed to start bridge for ~s: ~p",
[TunnelId, Reason]),
{error, Reason}
end.
%%%===================================================================
%%% Internal — Encryption
%%%===================================================================
tunnel_key() ->
Cookie = atom_to_binary(erlang:get_cookie()),
crypto:hash(sha256, <<"macula-dist-tunnel:", Cookie/binary>>).
%%%===================================================================
%%% Internal — Metrics
%%%===================================================================
init_metrics(TunnelId) ->
Ref = counters:new(4, [write_concurrency]),
Tunnels = persistent_term:get(macula_dist_tunnels, #{}),
persistent_term:put(macula_dist_tunnels, Tunnels#{TunnelId => Ref}),
Ref.
read_metrics(Ref) ->
#{bytes_out => counters:get(Ref, ?METRIC_BYTES_OUT),
bytes_in => counters:get(Ref, ?METRIC_BYTES_IN),
msgs_out => counters:get(Ref, ?METRIC_MSGS_OUT),
msgs_in => counters:get(Ref, ?METRIC_MSGS_IN)}.
%%%===================================================================
%%% Internal — Loopback Pair + Helpers
%%%===================================================================
create_loopback_pair() ->
ListenOpts = [binary, {active, false}, {reuseaddr, true}, {ip, {127,0,0,1}}],
{ok, LSock} = gen_tcp:listen(0, ListenOpts),
{ok, Port} = inet:port(LSock),
{ok, CSock} = gen_tcp:connect({127,0,0,1}, Port,
[binary, {active, false}, {packet, 2}, {nodelay, true}]),
{ok, ASock} = gen_tcp:accept(LSock),
gen_tcp:close(LSock),
inet:setopts(ASock, [{packet, raw}, {nodelay, true}]),
{CSock, ASock}.
flush_buffered_to_socket(BridgeSock, Key) ->
receive
{dist_data_buffered, EncData} ->
case decrypt(Key, EncData) of
{ok, Data} -> gen_tcp:send(BridgeSock, Data);
{error, _} -> gen_tcp:send(BridgeSock, EncData)
end,
flush_buffered_to_socket(BridgeSock, Key)
after 0 ->
ok
end.
decrypt(Key, <<Nonce:12/binary, Tag:16/binary, Ciphertext/binary>>) ->
case crypto:crypto_one_time_aead(
aes_256_gcm, Key, Nonce, Ciphertext, <<>>, Tag, false) of
error -> {error, decrypt_failed};
Plaintext -> {ok, Plaintext}
end;
decrypt(_Key, _Data) ->
{error, decrypt_failed}.
extract_payload(#{payload := P}) -> P;
extract_payload(P) when is_binary(P) -> P.