Packages

macula

4.2.3
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_dist_system macula_dist_relay.erl
Raw

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_pool/0]).
-export([register_mesh_pool/1]).
-export([advertise_dist_accept/0]).
-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).
%% Realm tag stamped on every dist tunnel frame. Dist tunnel traffic
%% is protocol-internal infrastructure (like the DHT) and not bound
%% to any user realm; the all-zeros tag is the SDK convention for
%% realm-agnostic infrastructure traffic.
-define(DIST_REALM, <<0:256>>).
-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 Macula V2 pool (`macula_client:pool()') as the
%% carrier for distribution tunnel traffic. Stored in `persistent_term'
%% so the bridge processes (which run inside dist_util's setup process,
%% not in any supervised tree) can pick it up.
-spec register_mesh_pool(pid()) -> ok.
register_mesh_pool(Pid) when is_pid(Pid) ->
persistent_term:put(macula_dist_mesh_pool, Pid),
?LOG_INFO("[dist_relay] Mesh pool 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_pool() of
undefined ->
{error, no_mesh_connection};
Pool ->
request_tunnel(Pool, 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_pool() of
undefined ->
?LOG_WARNING("[dist_relay] Cannot advertise — no mesh pool registered"),
ok;
Pool ->
NodeName = atom_to_binary(node()),
Procedure = <<"_dist.tunnel.", NodeName/binary>>,
Handler = fun(Args) -> handle_tunnel_request(Args) end,
macula_client:advertise(Pool, ?DIST_REALM, 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 Pool Lookup
%%%===================================================================
-spec get_mesh_pool() -> pid() | undefined.
get_mesh_pool() ->
pool_or_undef(persistent_term:get(macula_dist_mesh_pool, undefined)).
pool_or_undef(undefined) -> undefined;
pool_or_undef(Pid) when is_pid(Pid) ->
pool_alive(is_process_alive(Pid), Pid).
pool_alive(true, Pid) -> Pid;
pool_alive(false, _Pid) -> undefined.
%%%===================================================================
%%% Internal — Tunnel Negotiation (connecting side)
%%%===================================================================
request_tunnel(Pool, 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, Pool]),
%% V2 pool: first-success across healthy links. The pool itself
%% does the multi-station fan-out the V1 multi_relay used to do.
Result = macula_client:call(Pool, ?DIST_REALM, Procedure, Args,
?DIST_TIMEOUT),
on_tunnel_rpc_reply(Result, Pool).
on_tunnel_rpc_reply({ok, #{<<"tunnel_id">> := TunnelId}}, Pool) ->
?LOG_INFO("[dist_relay] Tunnel established: ~s", [TunnelId]),
create_dist_socket(Pool, TunnelId);
on_tunnel_rpc_reply({ok, #{<<"error">> := ErrorInfo}}, _Pool) ->
?LOG_WARNING("[dist_relay] Tunnel error: ~p", [ErrorInfo]),
{error, {tunnel_error, ErrorInfo}};
on_tunnel_rpc_reply({error, Reason}, _Pool) ->
?LOG_WARNING("[dist_relay] Tunnel request failed: ~p", [Reason]),
{error, {tunnel_failed, Reason}};
on_tunnel_rpc_reply(Other, _Pool) ->
?LOG_WARNING("[dist_relay] Unexpected RPC result: ~p", [Other]),
{error, {unexpected_result, Other}}.
%%%===================================================================
%%% 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">>,
on_tunnel_request_pool(get_mesh_pool(), TunnelId, SendTopic, RecvTopic).
on_tunnel_request_pool(undefined, _TunnelId, _SendTopic, _RecvTopic) ->
{error, <<"no_mesh_pool">>};
on_tunnel_request_pool(Pool, TunnelId, SendTopic, RecvTopic) ->
spawn_accept_bridge(Pool, TunnelId, SendTopic, RecvTopic),
{ok, #{<<"tunnel_id">> => TunnelId,
<<"send_topic">> => SendTopic,
<<"recv_topic">> => RecvTopic}}.
%% 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(Pool, TunnelId, SendTopic, RecvTopic) ->
{Pid, _Ref} = spawn_monitor(fun() ->
dist_accept_setup(Pool, 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(Pool, TunnelId, SendTopic, RecvTopic) ->
Self = self(),
Key = tunnel_key(),
%% Subscribe for buffering before kernel knows about us. The V2
%% pool delivers `{macula_event, SubRef, Topic, Payload, Meta}'
%% to Self directly; no callback indirection.
{ok, BufferSubRef} = macula_pubsub:subscribe(Pool, ?DIST_REALM,
RecvTopic, Self),
{DistSock, BridgeSock} = create_loopback_pair(),
on_kernel_lookup(whereis(net_kernel), Self, Pool, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId, Key).
on_kernel_lookup(undefined, _Self, Pool, _DistSock, _BridgeSock, BufferSubRef,
_SendTopic, _RecvTopic, _TunnelId, _Key) ->
?LOG_WARNING("[dist_bridge] net_kernel not found"),
macula_pubsub:unsubscribe(Pool, BufferSubRef);
on_kernel_lookup(KernelPid, Self, Pool, DistSock, BridgeSock, BufferSubRef,
SendTopic, RecvTopic, TunnelId, Key) ->
negotiate_with_kernel(Self, KernelPid, Pool, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId, Key).
negotiate_with_kernel(Self, KernelPid, Pool, 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_pubsub:unsubscribe(Pool, BufferSubRef),
flush_buffered_to_socket(BridgeSock, Key, BufferSubRef),
start_supervised_bridge(Pool, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key);
{KernelPid, unsupported_protocol} ->
?LOG_WARNING("[dist_bridge] Unsupported protocol"),
macula_pubsub:unsubscribe(Pool, BufferSubRef)
after ?CONTROLLER_TIMEOUT ->
?LOG_WARNING("[dist_bridge] Controller timeout"),
macula_pubsub:unsubscribe(Pool, BufferSubRef)
end.
%%%===================================================================
%%% Internal — Connect-side Bridge Creation
%%%===================================================================
create_dist_socket(Pool, TunnelId) ->
SendTopic = <<"_dist.data.", TunnelId/binary, ".out">>,
RecvTopic = <<"_dist.data.", TunnelId/binary, ".in">>,
{DistSock, BridgeSock} = create_loopback_pair(),
Key = tunnel_key(),
start_supervised_bridge(Pool, BridgeSock, SendTopic, RecvTopic,
TunnelId, Key),
{ok, DistSock, DistSock}.
%%%===================================================================
%%% Internal — Supervised Bridge Startup
%%%===================================================================
start_supervised_bridge(Pool, BridgeSock, SendTopic, RecvTopic, TunnelId, Key) ->
Metrics = init_metrics(TunnelId),
BridgeArgs = #{
pool => Pool,
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}.
%% Drain the buffered events that arrived while net_kernel was still
%% setting up the dist controller. The V2 pool delivers a tagged
%% 5-tuple — match on the SubRef we used to subscribe so we don't
%% accidentally consume unrelated events (the same process may host
%% other subscriptions during setup). `macula_event_gone' is the
%% terminal signal sent once the subscription tears down.
flush_buffered_to_socket(BridgeSock, Key, SubRef) ->
receive
{macula_event, SubRef, _Topic, EncData, _Meta} ->
ship_buffered(BridgeSock, decrypt(Key, EncData), EncData),
flush_buffered_to_socket(BridgeSock, Key, SubRef);
{macula_event_gone, SubRef, _Reason} ->
ok
after 0 ->
ok
end.
ship_buffered(BridgeSock, {ok, Data}, _EncData) ->
gen_tcp:send(BridgeSock, Data);
ship_buffered(BridgeSock, {error, _}, EncData) ->
gen_tcp:send(BridgeSock, EncData).
decrypt(Key, <<Nonce:12/binary, Tag:16/binary, Ciphertext/binary>>) ->
aead_decrypt(crypto:crypto_one_time_aead(
aes_256_gcm, Key, Nonce, Ciphertext, <<>>, Tag, false));
decrypt(_Key, _Data) ->
{error, decrypt_failed}.
aead_decrypt(error) -> {error, decrypt_failed};
aead_decrypt(Plaintext) -> {ok, Plaintext}.