Packages

macula

0.20.21
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_nat_system macula_hole_punch.erl
Raw

src/macula_nat_system/macula_hole_punch.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% QUIC Hole Punch Executor with Cancellation and Adaptive Timing.
%%%
%%% Implements the simultaneous open (SYN-SYN) pattern for QUIC
%%% to establish direct connections through NAT devices.
%%%
%%% Features:
%%% - Proper cancellation of in-progress punch attempts
%%% - Adaptive timing based on NAT type and previous attempts
%%% - Tracks active punches via gen_server state
%%% - Supports both sync and async execution
%%%
%%% QUIC Hole Punching Approach:
%%% Unlike TCP's explicit SYN packets, QUIC uses encrypted handshakes.
%%% The hole punching strategy is:
%%%
%%% 1. Both peers start QUIC connect() at the same coordinated time
%%% 2. Initial packets "punch" holes in both NATs
%%% 3. One peer's connection will succeed (race condition)
%%% 4. The other peer retries connecting through the opened hole
%%%
%%% NAT Behavior Considerations:
%%% - EI mapping: External address is consistent - easy hole punch
%%% - HD mapping: Must target specific host - coordinate addresses
%%% - PP allocation: Same port - single target port
%%% - PC allocation: Sequential ports - try predicted range
%%% - RD allocation: Random ports - harder to predict, try range
%%%
%%% Adaptive Timing:
%%% - Symmetric NAT: Longer timeouts, more port attempts
%%% - Restricted NAT: Standard timeouts
%%% - Full Cone: Fast timeouts, single port
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_hole_punch).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/0,
execute/3,
execute_async/3,
cancel/1,
get_active_punches/0
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-define(SERVER, ?MODULE).
%% Default timing parameters
-define(DEFAULT_PUNCH_TIMEOUT_MS, 2000).
-define(DEFAULT_CONNECT_TIMEOUT_MS, 1500).
-define(DEFAULT_MAX_PARALLEL, 3).
%% Adaptive timing parameters by NAT type
-define(SYMMETRIC_PUNCH_TIMEOUT_MS, 4000).
-define(SYMMETRIC_CONNECT_TIMEOUT_MS, 3000).
-define(SYMMETRIC_MAX_PARALLEL, 5).
-define(FULL_CONE_PUNCH_TIMEOUT_MS, 1000).
-define(FULL_CONE_CONNECT_TIMEOUT_MS, 800).
-define(FULL_CONE_MAX_PARALLEL, 1).
%%%===================================================================
%%% Types
%%%===================================================================
-type nat_type() :: full_cone | restricted | port_restricted | symmetric | unknown.
-type punch_opts() :: #{
target_host := binary() | string(),
target_ports := [inet:port_number()],
local_port => inet:port_number(),
session_id := binary(),
punch_time => integer(),
role => initiator | target,
local_nat_type => nat_type(),
remote_nat_type => nat_type()
}.
-type punch_result() ::
{ok, quicer:connection_handle()} |
{error, timeout | unreachable | all_ports_failed | cancelled}.
-record(state, {
%% Map from Ref -> {WorkerPid, [AttemptPids], ReplyTo}
active_punches = #{} :: #{reference() => {pid(), [pid()], pid()}}
}).
-export_type([punch_opts/0, punch_result/0, nat_type/0]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the hole punch executor.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [], []).
%% @doc Execute a hole punch attempt synchronously.
%% Blocks until connection established or timeout.
-spec execute(binary(), punch_opts(), timeout()) -> punch_result().
execute(TargetNodeId, Opts, Timeout) ->
gen_server:call(?SERVER, {execute, TargetNodeId, Opts, Timeout}, Timeout + 1000).
%% @doc Execute hole punch asynchronously.
%% Returns immediately, caller receives result via message.
-spec execute_async(binary(), punch_opts(), pid()) -> reference().
execute_async(TargetNodeId, Opts, ReplyTo) ->
gen_server:call(?SERVER, {execute_async, TargetNodeId, Opts, ReplyTo}).
%% @doc Cancel an ongoing hole punch attempt.
-spec cancel(reference()) -> ok | {error, not_found}.
cancel(Ref) ->
gen_server:call(?SERVER, {cancel, Ref}).
%% @doc Get list of active punch attempts (for debugging/monitoring).
-spec get_active_punches() -> [{reference(), map()}].
get_active_punches() ->
gen_server:call(?SERVER, get_active_punches).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init([]) ->
process_flag(trap_exit, true),
?LOG_INFO("Hole punch executor started"),
{ok, #state{}}.
handle_call({execute, TargetNodeId, Opts, Timeout}, From, State) ->
Ref = make_ref(),
{_ReplyPid, _} = From,
%% Spawn worker to do the actual punch
Self = self(),
WorkerPid = spawn_link(fun() ->
Result = do_execute(TargetNodeId, Opts, Timeout),
Self ! {punch_complete, Ref, Result}
end),
NewState = State#state{
active_punches = maps:put(Ref, {WorkerPid, [], From}, State#state.active_punches)
},
%% Don't reply now - reply when punch completes
{noreply, NewState};
handle_call({execute_async, TargetNodeId, Opts, ReplyTo}, _From, State) ->
Ref = make_ref(),
%% Spawn worker to do the actual punch
Self = self(),
WorkerPid = spawn_link(fun() ->
Timeout = get_adaptive_timeout(Opts),
Result = do_execute(TargetNodeId, Opts, Timeout),
Self ! {punch_complete, Ref, Result}
end),
NewState = State#state{
active_punches = maps:put(Ref, {WorkerPid, [], ReplyTo}, State#state.active_punches)
},
{reply, Ref, NewState};
handle_call({cancel, Ref}, _From, State) ->
case maps:get(Ref, State#state.active_punches, undefined) of
undefined ->
{reply, {error, not_found}, State};
{WorkerPid, AttemptPids, ReplyTo} ->
?LOG_DEBUG("Cancelling hole punch ~p with ~p attempts",
[Ref, length(AttemptPids)]),
%% Kill all attempt processes
lists:foreach(fun(Pid) ->
exit(Pid, cancelled)
end, [WorkerPid | AttemptPids]),
%% Notify waiting caller if sync
notify_cancelled(ReplyTo, Ref),
NewState = State#state{
active_punches = maps:remove(Ref, State#state.active_punches)
},
{reply, ok, NewState}
end;
handle_call(get_active_punches, _From, State) ->
Punches = maps:fold(fun(Ref, {WorkerPid, AttemptPids, _}, Acc) ->
[{Ref, #{
worker => WorkerPid,
attempts => length(AttemptPids),
alive => is_process_alive(WorkerPid)
}} | Acc]
end, [], State#state.active_punches),
{reply, Punches, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast({register_attempts, Ref, Pids}, State) ->
case maps:get(Ref, State#state.active_punches, undefined) of
undefined ->
{noreply, State};
{WorkerPid, _OldPids, ReplyTo} ->
NewState = State#state{
active_punches = maps:put(Ref, {WorkerPid, Pids, ReplyTo},
State#state.active_punches)
},
{noreply, NewState}
end;
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info({punch_complete, Ref, Result}, State) ->
case maps:get(Ref, State#state.active_punches, undefined) of
undefined ->
%% Already cancelled
{noreply, State};
{_WorkerPid, _AttemptPids, ReplyTo} ->
%% Notify the caller
notify_result(ReplyTo, Ref, Result),
NewState = State#state{
active_punches = maps:remove(Ref, State#state.active_punches)
},
{noreply, NewState}
end;
handle_info({'EXIT', Pid, Reason}, State) ->
?LOG_DEBUG("Hole punch process ~p exited: ~p", [Pid, Reason]),
%% Clean up any punches associated with this pid
NewPunches = maps:filter(fun(_Ref, {WorkerPid, _, _}) ->
WorkerPid =/= Pid
end, State#state.active_punches),
{noreply, State#state{active_punches = NewPunches}};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, State) ->
%% Cancel all active punches
maps:foreach(fun(_Ref, {WorkerPid, AttemptPids, _}) ->
lists:foreach(fun(Pid) -> exit(Pid, shutdown) end, [WorkerPid | AttemptPids])
end, State#state.active_punches),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private
%% @doc Notify caller of result (handles both sync and async)
notify_result({Pid, Tag}, _Ref, Result) when is_pid(Pid) ->
%% Sync call - reply via gen_server
gen_server:reply({Pid, Tag}, Result);
notify_result(Pid, Ref, Result) when is_pid(Pid) ->
%% Async call - send message
Pid ! {hole_punch_result, Ref, Result}.
%% @private
notify_cancelled({Pid, Tag}, _Ref) when is_pid(Pid) ->
gen_server:reply({Pid, Tag}, {error, cancelled});
notify_cancelled(Pid, Ref) when is_pid(Pid) ->
Pid ! {hole_punch_result, Ref, {error, cancelled}}.
%% @private
%% @doc Actual execution logic (runs in spawned process)
-spec do_execute(binary(), punch_opts(), timeout()) -> punch_result().
do_execute(TargetNodeId, Opts, Timeout) ->
#{target_host := Host, target_ports := Ports} = Opts,
?LOG_DEBUG("Executing hole punch to ~s (~s) on ports ~p",
[TargetNodeId, Host, Ports]),
%% Wait until coordinated punch time if specified
maybe_wait_for_punch_time(Opts),
%% Get adaptive parameters
{MaxParallel, ConnectTimeout} = get_adaptive_params(Opts),
%% Try connecting to target ports in parallel
attempt_connection(Host, Ports, Opts, Timeout, MaxParallel, ConnectTimeout).
%% @private
%% @doc Get adaptive timeout based on NAT types.
-spec get_adaptive_timeout(punch_opts()) -> timeout().
get_adaptive_timeout(Opts) ->
LocalNat = maps:get(local_nat_type, Opts, unknown),
RemoteNat = maps:get(remote_nat_type, Opts, unknown),
case worst_nat_type(LocalNat, RemoteNat) of
symmetric -> ?SYMMETRIC_PUNCH_TIMEOUT_MS;
full_cone -> ?FULL_CONE_PUNCH_TIMEOUT_MS;
_ -> ?DEFAULT_PUNCH_TIMEOUT_MS
end.
%% @private
%% @doc Get adaptive parameters based on NAT types.
-spec get_adaptive_params(punch_opts()) -> {pos_integer(), timeout()}.
get_adaptive_params(Opts) ->
LocalNat = maps:get(local_nat_type, Opts, unknown),
RemoteNat = maps:get(remote_nat_type, Opts, unknown),
case worst_nat_type(LocalNat, RemoteNat) of
symmetric ->
{?SYMMETRIC_MAX_PARALLEL, ?SYMMETRIC_CONNECT_TIMEOUT_MS};
full_cone ->
{?FULL_CONE_MAX_PARALLEL, ?FULL_CONE_CONNECT_TIMEOUT_MS};
_ ->
{?DEFAULT_MAX_PARALLEL, ?DEFAULT_CONNECT_TIMEOUT_MS}
end.
%% @private
%% @doc Determine the worst (most restrictive) NAT type.
-spec worst_nat_type(nat_type(), nat_type()) -> nat_type().
worst_nat_type(symmetric, _) -> symmetric;
worst_nat_type(_, symmetric) -> symmetric;
worst_nat_type(port_restricted, _) -> port_restricted;
worst_nat_type(_, port_restricted) -> port_restricted;
worst_nat_type(restricted, _) -> restricted;
worst_nat_type(_, restricted) -> restricted;
worst_nat_type(full_cone, full_cone) -> full_cone;
worst_nat_type(_, _) -> unknown.
%% @private
%% @doc Wait until the coordinated punch time arrives.
-spec maybe_wait_for_punch_time(punch_opts()) -> ok.
maybe_wait_for_punch_time(#{punch_time := PunchTime}) when is_integer(PunchTime) ->
Now = erlang:system_time(millisecond),
WaitMs = max(0, PunchTime - Now),
case WaitMs > 0 of
true ->
?LOG_DEBUG("Waiting ~p ms for coordinated punch time", [WaitMs]),
timer:sleep(WaitMs);
false ->
ok
end;
maybe_wait_for_punch_time(_) ->
ok.
%% @private
%% @doc Attempt connection to target on multiple ports.
-spec attempt_connection(binary() | string(), [inet:port_number()],
punch_opts(), timeout(), pos_integer(), timeout()) ->
punch_result().
attempt_connection(_Host, [], _Opts, _Timeout, _MaxPar, _ConnTimeout) ->
{error, all_ports_failed};
attempt_connection(Host, Ports, Opts, Timeout, MaxParallel, ConnectTimeout) ->
%% Limit parallel attempts
AttemptsCount = min(length(Ports), MaxParallel),
PortsToTry = lists:sublist(Ports, AttemptsCount),
?LOG_DEBUG("Trying ~p ports in parallel: ~p (timeout: ~p, conn_timeout: ~p)",
[AttemptsCount, PortsToTry, Timeout, ConnectTimeout]),
%% Spawn parallel connection attempts
Self = self(),
Ref = make_ref(),
Pids = [spawn_link(fun() ->
Result = try_single_connection(Host, Port, Opts, ConnectTimeout),
Self ! {punch_attempt, Ref, Port, Result}
end) || Port <- PortsToTry],
%% Wait for first success or all failures
collect_results(Ref, length(Pids), Timeout, Pids).
%% @private
%% @doc Try a single QUIC connection to host:port.
-spec try_single_connection(binary() | string(), inet:port_number(),
punch_opts(), timeout()) ->
{ok, quicer:connection_handle()} | {error, term()}.
try_single_connection(Host, Port, Opts, ConnectTimeout) ->
?LOG_DEBUG("Attempting QUIC connect to ~s:~p (timeout: ~p)",
[Host, Port, ConnectTimeout]),
%% Build QUIC connection options
ConnOpts = build_quic_opts(Opts, ConnectTimeout),
%% Convert host to appropriate format
HostStr = case is_binary(Host) of
true -> binary_to_list(Host);
false -> Host
end,
%% Attempt QUIC connection
case quicer:connect(HostStr, Port, ConnOpts, ConnectTimeout) of
{ok, Conn} ->
?LOG_INFO("Hole punch SUCCESS: connected to ~s:~p", [Host, Port]),
{ok, Conn};
{error, Reason} = Error ->
?LOG_DEBUG("Hole punch to ~s:~p failed: ~p", [Host, Port, Reason]),
Error
end.
%% @private
%% @doc Build QUIC connection options for hole punching.
-spec build_quic_opts(punch_opts(), timeout()) -> list().
build_quic_opts(Opts, ConnectTimeout) ->
BaseOpts = [
{alpn, ["macula"]},
{verify, none},
%% Adaptive timeouts based on NAT type
{idle_timeout_ms, ConnectTimeout * 3},
{handshake_idle_timeout_ms, ConnectTimeout},
%% Keep-alive to maintain NAT binding
{keep_alive_interval_ms, 1000}
],
%% Add local port binding if specified (for predictable source port)
case maps:get(local_port, Opts, undefined) of
undefined -> BaseOpts;
LocalPort -> [{local_port, LocalPort} | BaseOpts]
end.
%% @private
%% @doc Collect results from parallel connection attempts.
-spec collect_results(reference(), non_neg_integer(), timeout(), [pid()]) ->
punch_result().
collect_results(_Ref, 0, _Timeout, _Pids) ->
{error, all_ports_failed};
collect_results(Ref, Remaining, Timeout, Pids) ->
receive
{punch_attempt, Ref, Port, {ok, Conn}} ->
?LOG_DEBUG("Connection succeeded on port ~p, killing other attempts", [Port]),
%% Kill remaining attempts
lists:foreach(fun(Pid) -> exit(Pid, shutdown) end, Pids),
{ok, Conn};
{punch_attempt, Ref, Port, {error, Reason}} ->
?LOG_DEBUG("Connection to port ~p failed: ~p", [Port, Reason]),
collect_results(Ref, Remaining - 1, Timeout, Pids)
after Timeout ->
?LOG_DEBUG("Hole punch timeout waiting for ~p attempts", [Remaining]),
%% Kill all attempts on timeout
lists:foreach(fun(Pid) -> exit(Pid, shutdown) end, Pids),
{error, timeout}
end.