Packages
macula
0.42.0
7.1.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]).
-define(DIST_TIMEOUT, 25000).
-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() ->
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_relay_client:advertise(Client, Procedure, Handler),
?LOG_INFO("[dist_relay] Advertised distribution accept: ~s", [Procedure]),
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)},
%% Gap 7: try all connected relays when using multi_relay.
%% If the target is on a different relay, call_any tries each in sequence.
Result = tunnel_rpc(MeshClient, Procedure, Args),
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}}
end.
%% Try call_any first (multi_relay). Fall back to direct call (single relay_client).
tunnel_rpc(MeshClient, Procedure, Args) ->
case catch macula_multi_relay:call_any(MeshClient, Procedure, Args, ?DIST_TIMEOUT) of
{'EXIT', {noproc, _}} ->
%% Not a multi_relay — try as direct relay_client
macula_relay_client:call(MeshClient, Procedure, Args, ?DIST_TIMEOUT);
Result ->
Result
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_relay_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_relay_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},
macula_relay_client:unsubscribe(Client, BufferSubRef),
flush_buffered_to_socket(BridgeSock, Key),
%% Hand off to supervised bridge
start_supervised_bridge(Client, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key);
{KernelPid, unsupported_protocol} ->
?LOG_WARNING("[dist_bridge] Unsupported protocol"),
macula_relay_client:unsubscribe(Client, BufferSubRef)
after ?CONTROLLER_TIMEOUT ->
?LOG_WARNING("[dist_bridge] Controller timeout"),
macula_relay_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} ->
?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.