Packages
macula
0.41.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
%%%
%%% Tunnel bytes are encrypted with AES-256-GCM derived from the
%%% Erlang distribution cookie. The relay cannot read ETF content.
%%%
%%% This is EXPERIMENTAL. Supports Pid ! Msg, gen_server:call,
%%% pg groups, process monitoring. Does NOT guarantee Mnesia or
%%% global module compatibility over WAN latency.
%%%
%%% == 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.
%%% BRIDGE_RECV_TIMEOUT should be >= net_ticktime * 4.
%%% @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(BRIDGE_RECV_TIMEOUT, 60000).
-define(CONTROLLER_TIMEOUT, 30000).
-define(BACKPRESSURE_HWM, 64).
%% Metrics counter indices
-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.
%% Returns {ok, DistSock, DistSock} where DistSock is a gen_tcp loopback
%% socket bridged to an encrypted relay tunnel.
-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)},
case macula_relay_client:call(MeshClient, Procedure, Args, ?DIST_TIMEOUT) 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.
%%%===================================================================
%%% 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_monitor_bridge(Client, TunnelId, SendTopic, RecvTopic, FromNode),
{ok, #{<<"tunnel_id">> => TunnelId,
<<"send_topic">> => SendTopic,
<<"recv_topic">> => RecvTopic}}
end.
spawn_monitor_bridge(Client, TunnelId, SendTopic, RecvTopic, FromNode) ->
{Pid, _Ref} = spawn_monitor(fun() ->
dist_accept_bridge(Client, TunnelId, SendTopic, RecvTopic, FromNode)
end),
?LOG_INFO("[dist_relay] Tunnel bridge ~p for ~s", [Pid, TunnelId]),
Pid.
%%%===================================================================
%%% Internal — Accept-side Bridge
%%%===================================================================
dist_accept_bridge(Client, TunnelId, SendTopic, RecvTopic, _FromNode) ->
process_flag(trap_exit, true),
Self = self(),
Metrics = init_metrics(TunnelId),
Key = tunnel_key(),
{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");
KernelPid ->
setup_accept_bridge(Self, KernelPid, Client, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId,
Key, Metrics)
end.
setup_accept_bridge(Self, KernelPid, Client, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId,
Key, Metrics) ->
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),
start_bridge_io(Self, Client, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key, Metrics);
{KernelPid, unsupported_protocol} ->
?LOG_WARNING("[dist_bridge] Unsupported protocol")
after ?CONTROLLER_TIMEOUT ->
?LOG_WARNING("[dist_bridge] Controller timeout")
end.
start_bridge_io(Self, Client, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key, Metrics) ->
{ok, SubRef} = macula_relay_client:subscribe(Client, RecvTopic,
fun(Msg) -> Self ! {tunnel_in, extract_payload(Msg)} end),
spawn_link(fun() ->
bridge_reader_loop(Client, BridgeSock, SendTopic, TunnelId, Key, Metrics)
end),
bridge_writer_loop(BridgeSock, TunnelId, Client, SubRef, Key, Metrics).
%%%===================================================================
%%% Internal — Connect-side Bridge
%%%===================================================================
create_dist_socket(MeshClient, TunnelId) ->
SendTopic = <<"_dist.data.", TunnelId/binary, ".out">>,
RecvTopic = <<"_dist.data.", TunnelId/binary, ".in">>,
{DistSock, BridgeSock} = create_loopback_pair(),
spawn_link(fun() ->
tunnel_io_bridge(MeshClient, BridgeSock, SendTopic, RecvTopic, TunnelId)
end),
{ok, DistSock, DistSock}.
tunnel_io_bridge(MeshClient, BridgeSock, SendTopic, RecvTopic, TunnelId) ->
process_flag(trap_exit, true),
Self = self(),
Metrics = init_metrics(TunnelId),
Key = tunnel_key(),
{ok, SubRef} = macula_relay_client:subscribe(MeshClient, RecvTopic,
fun(Msg) -> Self ! {tunnel_in, extract_payload(Msg)} end),
spawn_link(fun() ->
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics)
end),
bridge_writer_loop(BridgeSock, TunnelId, MeshClient, SubRef, Key, Metrics).
%%%===================================================================
%%% Internal — Bridge I/O Loops
%%%===================================================================
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics) ->
case gen_tcp:recv(BridgeSock, 0, ?BRIDGE_RECV_TIMEOUT) of
{ok, Data} ->
maybe_backpressure(MeshClient),
Encrypted = encrypt(Key, Data),
macula_relay_client:publish(MeshClient, SendTopic, Encrypted),
counters:add(Metrics, ?METRIC_BYTES_OUT, byte_size(Data)),
counters:add(Metrics, ?METRIC_MSGS_OUT, 1),
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics);
{error, closed} ->
?LOG_INFO("[io_bridge] Reader closed for ~s", [TunnelId]);
{error, Reason} ->
?LOG_WARNING("[io_bridge] Reader error ~p for ~s", [Reason, TunnelId])
end.
%% Gap 1: cleanup on exit. Gap 8: writer timeout.
bridge_writer_loop(BridgeSock, TunnelId, Client, SubRef, Key, Metrics) ->
receive
{tunnel_in, EncData} when is_binary(EncData) ->
case decrypt(Key, EncData) of
{ok, Data} ->
counters:add(Metrics, ?METRIC_BYTES_IN, byte_size(Data)),
counters:add(Metrics, ?METRIC_MSGS_IN, 1),
case gen_tcp:send(BridgeSock, Data) of
ok ->
bridge_writer_loop(BridgeSock, TunnelId, Client,
SubRef, Key, Metrics);
{error, Reason} ->
?LOG_WARNING("[io_bridge] Writer send error ~p for ~s",
[Reason, TunnelId]),
cleanup_bridge(BridgeSock, TunnelId, Client, SubRef, Metrics)
end;
{error, decrypt_failed} ->
?LOG_WARNING("[io_bridge] Decrypt failed for ~s (cookie mismatch?)",
[TunnelId]),
bridge_writer_loop(BridgeSock, TunnelId, Client, SubRef, Key, Metrics)
end;
{'EXIT', _Pid, _Reason} ->
?LOG_INFO("[io_bridge] Linked process exited for ~s", [TunnelId]),
cleanup_bridge(BridgeSock, TunnelId, Client, SubRef, Metrics);
stop ->
cleanup_bridge(BridgeSock, TunnelId, Client, SubRef, Metrics)
after ?BRIDGE_RECV_TIMEOUT ->
?LOG_WARNING("[io_bridge] Writer timeout for ~s", [TunnelId]),
cleanup_bridge(BridgeSock, TunnelId, Client, SubRef, Metrics)
end.
%% Gap 1: proper cleanup — unsubscribe, close socket, remove metrics
cleanup_bridge(BridgeSock, TunnelId, Client, SubRef, _Metrics) ->
?LOG_INFO("[io_bridge] Cleaning up tunnel ~s", [TunnelId]),
catch macula_relay_client:unsubscribe(Client, SubRef),
catch gen_tcp:close(BridgeSock),
remove_metrics(TunnelId),
ok.
%%%===================================================================
%%% Internal — Encryption (Gap 3)
%%%===================================================================
%% Derive a symmetric key from the Erlang distribution cookie.
%% Both sides share the same cookie (required for dist handshake).
tunnel_key() ->
Cookie = atom_to_binary(erlang:get_cookie()),
crypto:hash(sha256, <<"macula-dist-tunnel:", Cookie/binary>>).
%% Encrypt with AES-256-GCM. Nonce is random per message (prepended).
encrypt(Key, Plaintext) ->
Nonce = crypto:strong_rand_bytes(12),
{Ciphertext, Tag} = crypto:crypto_one_time_aead(
aes_256_gcm, Key, Nonce, Plaintext, <<>>, true),
<<Nonce/binary, Tag/binary, Ciphertext/binary>>.
%% Decrypt AES-256-GCM. Nonce + tag prepended to ciphertext.
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}.
%%%===================================================================
%%% Internal — Backpressure (Gap 5)
%%%===================================================================
%% Pause if relay client's message queue is too deep.
maybe_backpressure(MeshClient) ->
case erlang:process_info(MeshClient, message_queue_len) of
{message_queue_len, Len} when Len > ?BACKPRESSURE_HWM ->
?LOG_WARNING("[io_bridge] Backpressure: relay queue ~p", [Len]),
timer:sleep(1);
_ ->
ok
end.
%%%===================================================================
%%% Internal — Metrics (Gap 9)
%%%===================================================================
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.
remove_metrics(TunnelId) ->
case persistent_term:get(macula_dist_tunnels, undefined) of
undefined -> ok;
Tunnels ->
persistent_term:put(macula_dist_tunnels, maps:remove(TunnelId, Tunnels))
end.
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, _} ->
%% During buffer phase, data may be unencrypted
%% (from the connecting side's first message before
%% encryption was set up on that side)
gen_tcp:send(BridgeSock, EncData)
end,
flush_buffered_to_socket(BridgeSock, Key)
after 0 ->
ok
end.
extract_payload(#{payload := P}) -> P;
extract_payload(P) when is_binary(P) -> P.