Packages

macula

0.40.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
macula src macula_dist_system macula_dist.erl
Raw

src/macula_dist_system/macula_dist.erl

%%%-------------------------------------------------------------------
%%% @doc QUIC Distribution Carrier for Erlang.
%%%
%%% This module implements the Erlang distribution carrier interface
%%% using QUIC transport via quicer. It replaces inet_tcp_dist to enable
%%% distributed Erlang over QUIC with:
%%%
%%% - Built-in TLS 1.3 encryption (mandatory)
%%% - NAT-friendly UDP-based transport
%%% - Connection migration support
%%% - Stream multiplexing for message priorities
%%% - Decentralized discovery (no EPMD required)
%%%
%%% == Usage ==
%%%
%%% Start the VM with:
%%% erl -proto_dist macula -no_epmd -start_epmd false
%%%
%%% Or in vm.args:
%%% -proto_dist macula
%%% -no_epmd
%%% -start_epmd false
%%% -macula_dist_port 4433
%%%
%%% Architecture: net_kernel - macula_dist - quicer - UDP/QUIC - remote node
%%%
%%% Node naming convention: port@ip (e.g., 4433@192.168.1.100)
%%%
%%% @copyright 2025 Macula.io Apache-2.0
%%% @end
%%%-------------------------------------------------------------------
-module(macula_dist).
%% Distribution carrier callbacks (required by OTP)
-export([
listen/1,
accept/1,
accept_connection/5,
setup/5,
close/1,
select/1,
childspecs/0
]).
%% Additional exports for integration
-export([
is_node_name/1,
splitname/1,
address/0
]).
%% Internal exports
-export([
acceptor_loop/2,
do_accept/3,
do_setup/6
]).
-include_lib("kernel/include/net_address.hrl").
-include_lib("kernel/include/dist.hrl").
-include_lib("kernel/include/dist_util.hrl").
-define(FAMILY, inet).
-define(DRIVER, macula_dist).
-define(DEFAULT_PORT, 4433).
-define(ALPN, "macula-dist").
-define(HANDSHAKE_TIMEOUT, 30000). % 30 seconds
-define(CONNECT_TIMEOUT, 15000). % 15 seconds (relay tunnel needs more time)
%% Stream IDs for different message types
-define(STREAM_CONTROL, 0).
-define(STREAM_DIST, 1).
%%%===================================================================
%%% Distribution Carrier Callbacks
%%%===================================================================
%% @doc Return child specifications for the distribution supervisor.
%% Called by net_sup during startup.
-spec childspecs() -> [supervisor:child_spec()].
childspecs() ->
%% The dist controller processes are supervised by net_kernel
%% We don't need additional supervised children
[].
%% @doc Listen for incoming distribution connections.
%% Returns a "listen handle" used by accept/1.
-spec listen(atom()) -> {ok, {term(), #net_address{}, 1..3}} | {error, term()}.
listen(NodeName) ->
Port = get_dist_port(NodeName),
case start_quic_listener(Port) of
{ok, ListenerHandle} ->
Address = make_address(NodeName, Port),
%% Return: {ListenerHandle, NetAddress, AddressFamily}
%% AddressFamily: 1=inet, 2=inet6, 3=local
{ok, {ListenerHandle, Address, 1}};
{error, Reason} ->
error_logger:error_msg("macula_dist: listen failed: ~p~n", [Reason]),
{error, Reason}
end.
%% @doc Accept incoming connections.
%% Called in a loop by net_kernel.
-spec accept(term()) -> pid().
accept(ListenerHandle) ->
spawn_link(?MODULE, acceptor_loop, [self(), ListenerHandle]).
%% @doc Accept a distribution connection.
%% This is called when a connection is being accepted from a remote node.
%% Socket can be {QuicConn, Stream} (direct QUIC) or a gen_tcp port (relay tunnel).
-spec accept_connection(pid(), term(), node(), term(), term()) -> pid().
accept_connection(AcceptPid, Socket, MyNode, Allowed, SetupTime) when is_port(Socket) ->
%% gen_tcp socket from relay tunnel loopback bridge
spawn_link(?MODULE, do_accept, [
{AcceptPid, Socket, Socket, MyNode, Allowed, SetupTime},
self(),
connection_id()
]);
accept_connection(AcceptPid, {QuicConn, Stream}, MyNode, Allowed, SetupTime) ->
spawn_link(?MODULE, do_accept, [
{AcceptPid, QuicConn, Stream, MyNode, Allowed, SetupTime},
self(),
connection_id()
]).
%% @doc Setup an outgoing distribution connection.
%% Called when this node wants to connect to another node.
%% self() here is net_kernel — pass it as Kernel to the spawned process.
-spec setup(node(), term(), atom(), term(), term()) -> pid().
setup(Node, Type, MyNode, LongOrShortNames, SetupTime) ->
Kernel = self(),
spawn_link(?MODULE, do_setup, [
Kernel, Node, Type, MyNode, LongOrShortNames, SetupTime
]).
%% @doc Close a distribution connection.
-spec close(term()) -> ok.
close(Socket) when is_port(Socket) ->
catch gen_tcp:close(Socket),
ok;
close({S, S}) when is_port(S) ->
catch gen_tcp:close(S),
ok;
close({QuicConn, _Stream}) ->
catch quicer:close_connection(QuicConn),
ok;
close(QuicConn) when is_reference(QuicConn) ->
catch quicer:close_connection(QuicConn),
ok;
close(_) ->
ok.
%% @doc Check if this module should handle distribution to the given node.
%% Returns true for any valid name@host or port@host format.
-spec select(atom()) -> boolean().
select(Node) ->
case splitname(Node) of
{_Port, _Host} -> true;
false -> false
end.
%%%===================================================================
%%% Name Handling
%%%===================================================================
%% @doc Check if a string is a valid node name.
-spec is_node_name(string()) -> boolean().
is_node_name(Name) ->
case splitname(Name) of
{_Port, _Host} -> true;
_ -> false
end.
%% @doc Split a node name into port and host.
%% Supports two formats:
%% - port@host (e.g., 4433@192.168.1.100) — explicit port
%% - name@host (e.g., test@nanode1.example.com) — uses default port
-spec splitname(atom() | string()) -> {integer(), string()} | false.
splitname(NodeName) when is_atom(NodeName) ->
splitname(atom_to_list(NodeName));
splitname(NodeName) when is_list(NodeName) ->
case string:tokens(NodeName, "@") of
[NameOrPort, Host] ->
case catch list_to_integer(NameOrPort) of
Port when is_integer(Port), Port > 0, Port < 65536 ->
{Port, Host};
_ ->
%% Standard name@host format — use default dist port
DefaultPort = application:get_env(kernel, macula_dist_port, ?DEFAULT_PORT),
{DefaultPort, Host}
end;
_ ->
false
end.
%% @doc Return address information for this distribution.
-spec address() -> #net_address{}.
address() ->
{ok, Host} = inet:gethostname(),
Port = application:get_env(kernel, macula_dist_port, ?DEFAULT_PORT),
make_address(Port, Host).
%%%===================================================================
%%% Internal Functions - Listener
%%%===================================================================
%% @private Start QUIC listener for distribution
start_quic_listener(Port) ->
%% Get or generate TLS certificates
{CertFile, KeyFile} = get_tls_certs(),
ListenerOpts = [
{certfile, CertFile},
{keyfile, KeyFile},
{alpn, [?ALPN]},
{idle_timeout_ms, 60000},
{peer_bidi_stream_count, 100},
{peer_unidi_stream_count, 100}
],
case quicer:listen(Port, ListenerOpts) of
{ok, Listener} ->
error_logger:info_msg("macula_dist: listening on UDP port ~p~n", [Port]),
{ok, Listener};
{error, Reason} ->
{error, Reason}
end.
%% @private Acceptor loop - waits for incoming connections
acceptor_loop(Kernel, Listener) ->
%% Accept with passive mode for recv() to work
%% Stream notifications still come via messages for async_accept_stream
case quicer:accept(Listener, #{active => false}) of
{ok, Conn} ->
error_logger:info_msg("macula_dist: accepted connection ~p~n", [Conn]),
case quicer:handshake(Conn) of
{ok, Conn} ->
error_logger:info_msg("macula_dist: handshake complete~n"),
%% Open or accept a bidirectional stream for distribution
case accept_dist_stream(Conn) of
{ok, Stream} ->
error_logger:info_msg("macula_dist: got stream ~p~n", [Stream]),
%% Notify kernel of new connection immediately
Kernel ! {accept, self(), {Conn, Stream}, ?FAMILY, ?DRIVER},
receive
{Kernel, controller, DistCtrl} ->
error_logger:info_msg("macula_dist: transferring to ~p~n", [DistCtrl]),
%% Use handoff_stream to forward any queued active messages
case quicer:handoff_stream(Stream, DistCtrl, #{}) of
ok ->
error_logger:info_msg("macula_dist: stream handoff ok~n"),
%% Transfer connection ownership
case quicer:controlling_process(Conn, DistCtrl) of
ok ->
error_logger:info_msg("macula_dist: conn transfer ok~n"),
DistCtrl ! {self(), controller, ok};
{error, ConnErr} ->
error_logger:warning_msg("macula_dist: conn transfer failed: ~p~n", [ConnErr]),
close({Conn, Stream})
end;
{error, StreamErr} ->
error_logger:warning_msg("macula_dist: stream handoff failed: ~p~n", [StreamErr]),
close({Conn, Stream})
end;
{Kernel, unsupported_protocol} ->
close({Conn, Stream})
end;
{error, StreamReason} ->
error_logger:warning_msg(
"macula_dist: stream accept failed: ~p~n",
[StreamReason]
),
quicer:close_connection(Conn)
end;
{error, HandshakeReason} ->
error_logger:warning_msg(
"macula_dist: handshake failed: ~p~n",
[HandshakeReason]
),
quicer:close_connection(Conn)
end;
{error, closed} ->
%% Listener was closed
exit(normal);
{error, AcceptReason} ->
error_logger:warning_msg(
"macula_dist: accept failed: ~p~n",
[AcceptReason]
)
end,
acceptor_loop(Kernel, Listener).
%% @private Accept distribution stream on connection
%% Use active mode so data arrives as messages (quicer recv doesn't work with passive)
accept_dist_stream(Conn) ->
error_logger:info_msg("macula_dist: waiting for stream on connection~n"),
%% Use active mode - data will arrive as {quic, Data, Stream, ...} messages
case quicer:accept_stream(Conn, #{active => true}, ?HANDSHAKE_TIMEOUT) of
{ok, Stream} ->
error_logger:info_msg("macula_dist: accepted stream: ~p~n", [Stream]),
{ok, Stream};
{error, timeout} ->
error_logger:warning_msg("macula_dist: stream accept timeout~n"),
{error, stream_timeout};
{error, Reason} ->
error_logger:warning_msg("macula_dist: stream accept failed: ~p~n", [Reason]),
{error, Reason}
end.
%%%===================================================================
%%% Internal Functions - Accept Connection
%%%===================================================================
%% @private Handle incoming connection setup
do_accept({AcceptPid, QuicConn, Stream, MyNode, Allowed, SetupTime}, Kernel, _ConnId) ->
Timer = dist_util:start_timer(SetupTime),
error_logger:info_msg("macula_dist: do_accept started, waiting for controller~n"),
%% Wait for control transfer from acceptor
receive
{AcceptPid, controller, ok} ->
error_logger:info_msg("macula_dist: do_accept got controller message~n"),
ok
after ?HANDSHAKE_TIMEOUT ->
error_logger:warning_msg("macula_dist: do_accept controller timeout~n"),
dist_util:shutdown(?MODULE, 280, control_transfer_timeout)
end,
%% Wait for handoff_done from the handoff_stream call (QUIC only, not gen_tcp)
case is_port(Stream) of
true ->
ok;
false ->
receive
{handoff_done, Stream, HandoffData} ->
error_logger:info_msg("macula_dist: do_accept got handoff_done: ~p~n", [HandoffData]),
ok
after 5000 ->
error_logger:warning_msg("macula_dist: do_accept handoff_done timeout~n"),
ok
end
end,
error_logger:info_msg("macula_dist: do_accept stream/conn ownership received~n"),
Socket = {QuicConn, Stream},
HSData = #hs_data{
kernel_pid = Kernel,
this_node = MyNode,
socket = Socket,
timer = Timer,
allowed = Allowed,
f_send = fun quic_send/2,
f_recv = fun quic_recv/3,
f_setopts_pre_nodeup = fun quic_setopts_pre_nodeup/1,
f_setopts_post_nodeup = fun quic_setopts_post_nodeup/1,
f_getll = fun quic_getll/1,
f_address = fun quic_address/2,
f_handshake_complete = fun quic_handshake_complete/3,
mf_tick = fun quic_tick/1,
mf_getstat = fun quic_getstat/1,
mf_setopts = fun quic_setopts/2,
mf_getopts = fun quic_getopts/2
},
error_logger:info_msg("macula_dist: do_accept starting handshake~n"),
%% Run distribution handshake — catch any exit for diagnostics
try
dist_util:handshake_other_started(HSData)
catch
Class:Reason:Stack ->
error_logger:error_msg("macula_dist: handshake FAILED ~p:~p~n~p~n",
[Class, Reason, Stack])
end.
%%%===================================================================
%%% Internal Functions - Setup Outgoing Connection
%%%===================================================================
%% @private Setup outgoing distribution connection.
%% Two modes:
%% - Direct QUIC (default): connect_quic(Host, Port) — requires mutual reachability
%% - Relay mesh (MACULA_DIST_MODE=relay): tunnel through relay — outbound-only nodes
do_setup(Kernel, Node, Type, MyNode, _LongOrShortNames, SetupTime) ->
Timer = dist_util:start_timer(SetupTime),
case splitname(Node) of
{Port, Host} ->
ConnectFn = case macula_dist_relay:is_relay_mode() of
true -> fun() -> macula_dist_relay:connect(atom_to_list(Node), Host, Port) end;
false -> fun() -> connect_quic(Host, Port) end
end,
case ConnectFn() of
{ok, Conn, Stream} ->
setup_loop(Kernel, {Conn, Stream}, Node, Type, MyNode, Timer);
{error, Reason} ->
error_logger:warning_msg(
"macula_dist: connection to ~p failed: ~p~n",
[Node, Reason]
),
dist_util:shutdown(?MODULE, 318, {connect_failed, Reason})
end;
false ->
error_logger:warning_msg(
"macula_dist: invalid node name format: ~p~n",
[Node]
),
dist_util:shutdown(?MODULE, 325, invalid_node_name)
end.
%% @private Connect to remote node via QUIC
%% Uses macula_tls for certificate verification settings (v0.11.0+)
connect_quic(Host, Port) ->
{CertFile, KeyFile} = get_tls_certs(),
%% Get TLS verification options from centralized module
TlsOpts = macula_tls:quic_client_opts(),
%% Merge with distribution-specific options
BaseOpts = [
{alpn, [?ALPN]},
{certfile, CertFile},
{keyfile, KeyFile},
{idle_timeout_ms, 60000}
],
ConnOpts = merge_dist_opts(BaseOpts, TlsOpts),
error_logger:info_msg("macula_dist: connecting to ~s:~p~n", [Host, Port]),
case quicer:connect(Host, Port, ConnOpts, ?CONNECT_TIMEOUT) of
{ok, Conn} ->
error_logger:info_msg("macula_dist: connected, opening stream~n"),
%% Open bidirectional stream for distribution - use passive mode
case quicer:start_stream(Conn, #{active => false}) of
{ok, Stream} ->
error_logger:info_msg("macula_dist: stream opened: ~p~n", [Stream]),
{ok, Conn, Stream};
{error, StreamReason} ->
error_logger:warning_msg("macula_dist: stream open failed: ~p~n", [StreamReason]),
quicer:close_connection(Conn),
{error, {stream_failed, StreamReason}}
end;
{error, Reason} ->
error_logger:warning_msg("macula_dist: connect failed: ~p~n", [Reason]),
{error, Reason}
end.
%% @private Complete connection setup — Socket is either {QuicConn, QuicStream}
%% or a gen_tcp socket (from relay tunnel loopback pair).
setup_loop(Kernel, Socket, Node, Type, MyNode, Timer) ->
HSData = #hs_data{
kernel_pid = Kernel,
this_node = MyNode,
other_node = Node,
socket = Socket,
timer = Timer,
this_flags = 0,
other_flags = 0,
f_send = fun quic_send/2,
f_recv = fun quic_recv/3,
f_setopts_pre_nodeup = fun quic_setopts_pre_nodeup/1,
f_setopts_post_nodeup = fun quic_setopts_post_nodeup/1,
f_getll = fun quic_getll/1,
f_address = fun quic_address/2,
f_handshake_complete = fun quic_handshake_complete/3,
mf_tick = fun quic_tick/1,
mf_getstat = fun quic_getstat/1,
mf_setopts = fun quic_setopts/2,
mf_getopts = fun quic_getopts/2,
request_type = Type
},
try
dist_util:handshake_we_started(HSData)
catch
Class:Reason:Stack ->
error_logger:error_msg("macula_dist: handshake_we_started FAILED ~p:~p~n~p~n",
[Class, Reason, Stack])
end.
%%%===================================================================
%%% QUIC Socket Operations
%%%===================================================================
%% @private Send data over QUIC stream or relay tunnel
%% gen_tcp socket (from relay tunnel loopback pair)
quic_send(Socket, Data) when is_port(Socket) ->
gen_tcp:send(Socket, Data);
%% Tuple socket: {Socket, Socket} from tunnel or {Conn, Stream} from QUIC
quic_send({S, S}, Data) when is_port(S) ->
gen_tcp:send(S, Data);
quic_send({_Conn, Stream}, Data) ->
case quicer:send(Stream, Data) of
{ok, _} -> ok;
{error, Reason} ->
error_logger:warning_msg("macula_dist: quic_send error: ~p~n", [Reason]),
{error, Reason}
end.
%% @private Receive data from QUIC stream or relay tunnel
%% gen_tcp socket (from relay tunnel loopback pair)
%% dist_util expects {ok, List} not {ok, Binary} — convert
quic_recv(Socket, Length, Timeout) when is_port(Socket) ->
case gen_tcp:recv(Socket, Length, recv_timeout(Timeout)) of
{ok, Bin} when is_binary(Bin) -> {ok, binary_to_list(Bin)};
Other -> Other
end;
quic_recv({S, S}, Length, Timeout) when is_port(S) ->
case gen_tcp:recv(S, Length, recv_timeout(Timeout)) of
{ok, Bin} when is_binary(Bin) -> {ok, binary_to_list(Bin)};
Other -> Other
end;
quic_recv({_Conn, Stream}, _Length, Timeout) ->
TimeoutMs = recv_timeout(Timeout),
receive
{quic, Data, Stream, _Flags} when is_binary(Data) ->
{ok, binary_to_list(Data)};
{quic, stream_closed, Stream, _Flags} ->
{error, closed};
{quic, peer_send_shutdown, Stream, _Flags} ->
{error, closed}
after TimeoutMs ->
{error, timeout}
end.
recv_timeout(infinity) -> 30000;
recv_timeout(T) -> T.
%% Match inet_tcp_dist: {packet,4} for post-handshake distribution protocol
quic_setopts_pre_nodeup(Socket) when is_port(Socket) ->
inet:setopts(Socket, [{active, false}, {packet, 4}]);
quic_setopts_pre_nodeup({S, S}) when is_port(S) ->
inet:setopts(S, [{active, false}, {packet, 4}]);
quic_setopts_pre_nodeup({_Conn, _Stream}) ->
ok.
quic_setopts_post_nodeup(Socket) when is_port(Socket) ->
inet:setopts(Socket, [{active, true}, {packet, 4}, {deliver, port}, binary]);
quic_setopts_post_nodeup({S, S}) when is_port(S) ->
inet:setopts(S, [{active, true}, {packet, 4}, {deliver, port}, binary]);
quic_setopts_post_nodeup({_Conn, _Stream}) ->
ok.
quic_getll(Socket) when is_port(Socket) ->
{ok, Socket};
quic_getll({S, S}) when is_port(S) ->
{ok, S};
quic_getll({Conn, _Stream}) ->
{ok, Conn}.
%% @private Get address information
quic_address(Socket, Node) when is_port(Socket) ->
quic_address_for_tcp(Socket, Node);
quic_address({S, S}, Node) when is_port(S) ->
quic_address_for_tcp(S, Node);
quic_address({Conn, _Stream}, Node) ->
case quicer:peername(Conn) of
{ok, {IP, Port}} ->
#net_address{
address = {IP, Port},
host = atom_to_list(Node),
protocol = ?DRIVER,
family = ?FAMILY
};
{error, _} ->
#net_address{
address = undefined,
host = atom_to_list(Node),
protocol = ?DRIVER,
family = ?FAMILY
}
end.
quic_address_for_tcp(Sock, Node) ->
case inet:peername(Sock) of
{ok, {IP, Port}} ->
#net_address{
address = {IP, Port},
host = atom_to_list(Node),
protocol = ?DRIVER,
family = ?FAMILY
};
{error, _} ->
#net_address{
address = undefined,
host = atom_to_list(Node),
protocol = ?DRIVER,
family = ?FAMILY
}
end.
%% @private Called when distribution handshake completes
quic_handshake_complete(_Socket, _Node, _DHandle) ->
ok.
%% @private Send tick (keepalive)
quic_tick(Socket) when is_port(Socket) ->
gen_tcp:send(Socket, <<>>);
quic_tick({S, S}) when is_port(S) ->
gen_tcp:send(S, <<>>);
quic_tick({_Conn, Stream}) ->
quicer:send(Stream, <<>>).
%% @private Get socket statistics
%% dist_util expects {ok, RecvCount, SendCount, SendPend} — NOT a proplist
quic_getstat(Socket) when is_port(Socket) ->
getstat_tcp(Socket);
quic_getstat({S, S}) when is_port(S) ->
getstat_tcp(S);
quic_getstat({Conn, _Stream}) ->
case quicer:getstat(Conn, [send_cnt, recv_cnt, send_pend]) of
{ok, Stats} -> split_stat(Stats, 0, 0, 0);
{error, _} -> {ok, 0, 0, 0}
end.
getstat_tcp(Socket) ->
case inet:getstat(Socket, [recv_cnt, send_cnt, send_pend]) of
{ok, Stats} -> split_stat(Stats, 0, 0, 0);
Error -> Error
end.
split_stat([{recv_cnt, R} | Rest], _, W, P) -> split_stat(Rest, R, W, P);
split_stat([{send_cnt, W} | Rest], R, _, P) -> split_stat(Rest, R, W, P);
split_stat([{send_pend, P} | Rest], R, W, _) -> split_stat(Rest, R, W, P);
split_stat([], R, W, P) -> {ok, R, W, P}.
%% @private Set socket options
quic_setopts(Socket, Opts) when is_port(Socket) ->
inet:setopts(Socket, Opts);
quic_setopts({S, S}, Opts) when is_port(S) ->
inet:setopts(S, Opts);
quic_setopts({_Conn, Stream}, Opts) ->
lists:foreach(
fun({active, Mode}) -> quicer:setopt(Stream, active, Mode);
(_) -> ok
end,
Opts
),
ok.
%% @private Get socket options
quic_getopts(Socket, Opts) when is_port(Socket) ->
inet:getopts(Socket, Opts);
quic_getopts({S, S}, Opts) when is_port(S) ->
inet:getopts(S, Opts);
quic_getopts({_Conn, _Stream}, _Opts) ->
%% Return defaults for now
{ok, [{active, true}]}.
%%%===================================================================
%%% Utility Functions
%%%===================================================================
%% @private Merge distribution connection options
%% TLS options from macula_tls take precedence over base options
merge_dist_opts(BaseOpts, TlsOpts) ->
lists:foldl(
fun({Key, Value}, Acc) ->
lists:keystore(Key, 1, Acc, {Key, Value})
end,
BaseOpts,
TlsOpts
).
%% @private Get distribution port from node name or config
%% NodeName can be:
%% - Full name: '4433@192.168.1.100' (atom with port@host)
%% - Short name: '4433' (atom with just port - this is what listen/1 receives)
%% - String version of either
get_dist_port(NodeName) when is_atom(NodeName) ->
get_dist_port(atom_to_list(NodeName));
get_dist_port(NodeName) when is_list(NodeName) ->
%% First try full name format (port@host)
case splitname(NodeName) of
{Port, _Host} ->
Port;
false ->
%% Try parsing as just a port number
case catch list_to_integer(NodeName) of
Port when is_integer(Port), Port > 0, Port < 65536 ->
Port;
_ ->
application:get_env(kernel, macula_dist_port, ?DEFAULT_PORT)
end
end.
%% @private Get TLS certificates
get_tls_certs() ->
CertDir = application:get_env(kernel, macula_dist_cert_dir, "/tmp/macula_dist"),
CertFile = filename:join(CertDir, "cert.pem"),
KeyFile = filename:join(CertDir, "key.pem"),
%% Generate self-signed certs if they don't exist
case filelib:is_regular(CertFile) andalso filelib:is_regular(KeyFile) of
true ->
{CertFile, KeyFile};
false ->
generate_self_signed_cert(CertDir, CertFile, KeyFile)
end.
%% @private Generate self-signed certificate
generate_self_signed_cert(_CertDir, CertFile, KeyFile) ->
ok = filelib:ensure_dir(CertFile),
%% Use OpenSSL to generate self-signed cert
Cmd = io_lib:format(
"openssl req -x509 -newkey rsa:2048 -keyout ~s -out ~s "
"-days 365 -nodes -subj '/CN=macula-dist' 2>/dev/null",
[KeyFile, CertFile]
),
case os:cmd(lists:flatten(Cmd)) of
"" ->
{CertFile, KeyFile};
Error ->
error_logger:error_msg("Failed to generate certs: ~s~n", [Error]),
%% Return paths anyway - will fail later with clear error
{CertFile, KeyFile}
end.
%% @private Create net_address record
make_address(NodeName, Port) when is_atom(NodeName) ->
{ok, Host} = inet:gethostname(),
make_address(Port, Host);
make_address(Port, Host) when is_integer(Port) ->
#net_address{
address = {Host, Port},
host = Host,
protocol = ?DRIVER,
family = ?FAMILY
}.
%% @private Generate unique connection ID
connection_id() ->
erlang:unique_integer([positive, monotonic]).