Packages
macula
0.32.2
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_nat_system/macula_connection_upgrade.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Connection Upgrade Manager.
%%%
%%% Handles upgrading connections from relay (indirect) to direct
%%% when hole punching succeeds. This is a key optimization for
%%% reducing latency and bootstrap load.
%%%
%%% Upgrade Flow:
%%% 1. Connection established via relay (fallback)
%%% 2. Hole punch attempt happens in background
%%% 3. On success, upgrade_to_direct/2 is called
%%% 4. Messages are seamlessly transitioned to direct connection
%%% 5. Relay connection is gracefully closed
%%%
%%% Key Features:
%%% - Message ordering preserved during upgrade
%%% - No message loss during transition
%%% - Automatic fallback if direct connection fails
%%% - Metrics tracking for upgrade success/failure
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_connection_upgrade).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/0,
upgrade_to_direct/3,
register_relay/3,
unregister_relay/1,
get_upgrade_stats/0
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-define(SERVER, ?MODULE).
%% Grace period to drain in-flight messages before closing relay
-define(DRAIN_TIMEOUT_MS, 1000).
%%%===================================================================
%%% Types
%%%===================================================================
-record(relay_info, {
peer_id :: binary(),
relay_conn :: term(), % Connection handle to relay
endpoint :: {inet:ip_address(), inet:port_number()},
created_at :: integer()
}).
-record(state, {
%% Map from peer_id -> relay_info
relays = #{} :: #{binary() => #relay_info{}},
%% Upgrade statistics
stats = #{
upgrades_attempted => 0,
upgrades_succeeded => 0,
upgrades_failed => 0
} :: map()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the connection upgrade manager.
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [], []).
%% @doc Upgrade a relay connection to direct.
%% PeerId - The remote peer's ID
%% RelayConn - Current relay connection handle
%% DirectConn - New direct connection handle from hole punch
-spec upgrade_to_direct(binary(), term(), term()) ->
ok | {error, no_relay | upgrade_failed}.
upgrade_to_direct(PeerId, RelayConn, DirectConn) ->
gen_server:call(?SERVER, {upgrade, PeerId, RelayConn, DirectConn}, 10000).
%% @doc Register a relay connection for potential future upgrade.
%% Called when connection is established via relay fallback.
-spec register_relay(binary(), term(), {inet:ip_address(), inet:port_number()}) -> ok.
register_relay(PeerId, RelayConn, Endpoint) ->
gen_server:cast(?SERVER, {register_relay, PeerId, RelayConn, Endpoint}).
%% @doc Unregister a relay connection (closed or upgraded).
-spec unregister_relay(binary()) -> ok.
unregister_relay(PeerId) ->
gen_server:cast(?SERVER, {unregister_relay, PeerId}).
%% @doc Get upgrade statistics.
-spec get_upgrade_stats() -> map().
get_upgrade_stats() ->
gen_server:call(?SERVER, get_stats).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init([]) ->
?LOG_INFO("Connection upgrade manager started"),
{ok, #state{}}.
handle_call({upgrade, PeerId, RelayConn, DirectConn}, _From, State) ->
NewStats = increment_stat(upgrades_attempted, State#state.stats),
case maps:get(PeerId, State#state.relays, undefined) of
undefined ->
?LOG_WARNING("Upgrade requested for unknown relay: ~p", [PeerId]),
{reply, {error, no_relay}, State#state{stats = NewStats}};
#relay_info{relay_conn = RegisteredRelay} = RelayInfo ->
%% Verify the relay connection matches
case verify_relay(RelayConn, RegisteredRelay) of
true ->
%% Perform the upgrade
case do_upgrade(PeerId, RelayInfo, DirectConn) of
ok ->
?LOG_INFO("Connection upgraded to direct for peer ~p", [PeerId]),
FinalStats = increment_stat(upgrades_succeeded, NewStats),
NewRelays = maps:remove(PeerId, State#state.relays),
{reply, ok, State#state{
relays = NewRelays,
stats = FinalStats
}};
{error, Reason} ->
?LOG_WARNING("Upgrade failed for peer ~p: ~p", [PeerId, Reason]),
FinalStats = increment_stat(upgrades_failed, NewStats),
{reply, {error, upgrade_failed}, State#state{stats = FinalStats}}
end;
false ->
?LOG_WARNING("Relay connection mismatch for peer ~p", [PeerId]),
{reply, {error, relay_mismatch}, State#state{stats = NewStats}}
end
end;
handle_call(get_stats, _From, State) ->
BaseStats = State#state.stats,
Stats = BaseStats#{active_relays => maps:size(State#state.relays)},
{reply, Stats, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast({register_relay, PeerId, RelayConn, Endpoint}, State) ->
RelayInfo = #relay_info{
peer_id = PeerId,
relay_conn = RelayConn,
endpoint = Endpoint,
created_at = erlang:system_time(millisecond)
},
?LOG_DEBUG("Registered relay for peer ~p", [PeerId]),
NewRelays = maps:put(PeerId, RelayInfo, State#state.relays),
{noreply, State#state{relays = NewRelays}};
handle_cast({unregister_relay, PeerId}, State) ->
?LOG_DEBUG("Unregistered relay for peer ~p", [PeerId]),
NewRelays = maps:remove(PeerId, State#state.relays),
{noreply, State#state{relays = NewRelays}};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private
%% @doc Verify the relay connection is the one we registered.
verify_relay(ProvidedConn, RegisteredConn) ->
ProvidedConn =:= RegisteredConn.
%% @private
%% @doc Perform the actual connection upgrade.
%% Steps:
%% 1. Pause new messages on relay (mark as upgrading)
%% 2. Drain in-flight messages (short wait)
%% 3. Update routing table to use direct connection
%% 4. Close relay connection gracefully
-spec do_upgrade(binary(), #relay_info{}, term()) -> ok | {error, term()}.
do_upgrade(PeerId, RelayInfo, DirectConn) ->
handle_upgrade_result(catch execute_upgrade_steps(PeerId, RelayInfo, DirectConn)).
%% @private Execute upgrade steps in sequence
execute_upgrade_steps(PeerId, RelayInfo, DirectConn) ->
%% Step 1: Mark relay as upgrading (stop accepting new messages)
mark_relay_upgrading(RelayInfo),
%% Step 2: Short drain period for in-flight messages
timer:sleep(?DRAIN_TIMEOUT_MS),
%% Step 3: Update routing to use direct connection
update_routing(PeerId, DirectConn),
%% Step 4: Close relay connection gracefully
close_relay_gracefully(RelayInfo),
ok.
%% @private Handle upgrade result
handle_upgrade_result({'EXIT', {Reason, Stack}}) ->
?LOG_ERROR("Upgrade failed: ~p~n~p", [Reason, Stack]),
{error, Reason};
handle_upgrade_result({'EXIT', Reason}) ->
?LOG_ERROR("Upgrade failed: ~p", [Reason]),
{error, Reason};
handle_upgrade_result(ok) ->
ok;
handle_upgrade_result({error, _Reason} = Error) ->
Error.
%% @private
%% @doc Mark relay connection as upgrading (no new messages).
mark_relay_upgrading(#relay_info{relay_conn = Conn}) ->
%% In a full implementation, this would signal the connection
%% to stop accepting new outbound messages
?LOG_DEBUG("Marking relay ~p as upgrading", [Conn]),
ok.
%% @private
%% @doc Update routing table to use direct connection.
update_routing(PeerId, DirectConn) ->
%% Notify the connection pool about the new direct connection
case whereis(macula_peer_connection_pool) of
undefined ->
?LOG_DEBUG("Connection pool not running, routing update skipped");
_Pid ->
%% Store the direct connection in the pool
macula_peer_connection_pool:put(PeerId, DirectConn)
end,
%% Also update the peer system if running
case whereis(macula_peer_system) of
undefined ->
ok;
_PeerPid ->
%% The peer system may need to know about direct connections
?LOG_DEBUG("Notifying peer system of direct connection for ~p", [PeerId])
end,
ok.
%% @private
%% @doc Close relay connection gracefully.
close_relay_gracefully(#relay_info{relay_conn = Conn}) ->
%% Attempt graceful close, ignore errors
?LOG_DEBUG("Closing relay connection ~p", [Conn]),
case catch close_connection(Conn) of
ok -> ok;
{'EXIT', _} -> ok;
{error, _} -> ok
end.
%% @private
%% @doc Close a connection (abstraction for different connection types).
close_connection(Conn) when is_pid(Conn) ->
gen_server:stop(Conn, normal, 1000);
close_connection(_Conn) ->
%% QUIC connection handle - would use quicer:close/1
%% For now, just return ok
ok.
%% @private
%% @doc Increment a statistics counter.
increment_stat(Key, Stats) ->
maps:update_with(Key, fun(V) -> V + 1 end, 1, Stats).