Packages

macula

0.40.1
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
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
%%%
%%% This is EXPERIMENTAL. Supports Pid ! Msg, gen_server:call,
%%% pg groups, process monitoring. Does NOT guarantee Mnesia or
%%% global module compatibility over WAN latency.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_dist_relay).
-include_lib("kernel/include/logger.hrl").
-export([connect/3, accept_dist/2]).
-export([is_relay_mode/0, get_mesh_client/0]).
-export([register_mesh_client/1]).
-export([advertise_dist_accept/0]).
-define(DIST_TIMEOUT, 25000).
-define(BRIDGE_RECV_TIMEOUT, 60000).
-define(CONTROLLER_TIMEOUT, 30000).
%%%===================================================================
%%% Public API
%%%===================================================================
%% @doc Register a mesh relay client for distribution tunneling.
%% Called by the application that owns the relay client (e.g., hecate_mesh).
-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 a 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 Accept an incoming distribution connection via relay mesh.
-spec accept_dist(binary(), map()) -> {ok, binary()}.
accept_dist(TunnelId, _Opts) ->
?LOG_INFO("[dist_relay] Accepting dist tunnel: ~s", [TunnelId]),
{ok, TunnelId}.
%% @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.
%%%===================================================================
%%% Internal — Mesh Client Lookup
%%%===================================================================
%% @doc Find the mesh relay client from persistent_term.
-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
%%%===================================================================
%% Creates a gen_tcp loopback pair so OTP's dist_util gets a real fd.
%% One end goes to dist_util, the other is bridged to the relay tunnel.
dist_accept_bridge(Client, TunnelId, SendTopic, RecvTopic, _FromNode) ->
Self = self(),
%% Subscribe BEFORE notifying kernel — buffer data until bridge socket ready.
{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)
end.
setup_accept_bridge(Self, KernelPid, Client, DistSock, BridgeSock,
BufferSubRef, SendTopic, RecvTopic, TunnelId) ->
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),
start_bridge_io(Self, Client, BridgeSock, SendTopic, RecvTopic, TunnelId);
{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) ->
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) end),
bridge_writer_loop(BridgeSock, TunnelId).
%%%===================================================================
%%% 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) ->
Self = self(),
{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) end),
bridge_writer_loop(BridgeSock, TunnelId).
%%%===================================================================
%%% Internal — Bridge I/O Loops
%%%===================================================================
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId) ->
case gen_tcp:recv(BridgeSock, 0, ?BRIDGE_RECV_TIMEOUT) of
{ok, Data} ->
macula_relay_client:publish(MeshClient, SendTopic, Data),
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId);
{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.
bridge_writer_loop(BridgeSock, TunnelId) ->
receive
{tunnel_in, Data} when is_binary(Data) ->
case gen_tcp:send(BridgeSock, Data) of
ok -> bridge_writer_loop(BridgeSock, TunnelId);
{error, Reason} ->
?LOG_WARNING("[io_bridge] Writer error ~p for ~s", [Reason, TunnelId])
end;
stop ->
gen_tcp:close(BridgeSock)
end.
%%%===================================================================
%%% Internal — Loopback Pair + Helpers
%%%===================================================================
create_loopback_pair() ->
%% DistSock (CSock) gets {packet, 2} for handshake framing.
%% BridgeSock (ASock) gets {packet, raw} — transparent byte pipe.
%% When dist_util switches DistSock to {packet, 4} post-handshake,
%% the bridge doesn't care — it forwards raw bytes including headers.
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) ->
receive
{dist_data_buffered, Data} ->
gen_tcp:send(BridgeSock, Data),
flush_buffered_to_socket(BridgeSock)
after 0 ->
ok
end.
extract_payload(#{payload := P}) -> P;
extract_payload(P) when is_binary(P) -> P.