Packages
macula
0.35.1
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_connection_pool.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Connection pool manager for endpoint connections.
%%%
%%% Manages a pool of QUIC connections to remote endpoints, providing
%%% connection caching and reuse to avoid connection overhead for
%%% multi-endpoint RPC operations.
%%%
%%% Connection pool structure:
%%% #{Endpoint => #{connection => Conn, stream => Stream, last_used => Timestamp}}
%%% @end
%%%-------------------------------------------------------------------
-module(macula_connection_pool).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
%% API
-export([
get_or_create_connection/4,
create_connection/4,
close_all_connections/1
]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Get or create a connection to an endpoint.
%% Returns {ok, Conn, Stream, UpdatedPool} or {error, Reason, Pool}.
-spec get_or_create_connection(binary(), binary(), binary(), map()) ->
{ok, pid(), pid(), map()} | {error, term(), map()}.
get_or_create_connection(Endpoint, NodeId, RealmId, Pool) when is_map(Pool) ->
case maps:get(Endpoint, Pool, undefined) of
undefined ->
%% No existing connection, create new one
case create_connection(Endpoint, NodeId, RealmId, Pool) of
{ok, Conn, Stream, _Pool2} ->
%% Cache the connection
ConnectionInfo = #{
connection => Conn,
stream => Stream,
last_used => erlang:system_time(second)
},
UpdatedPool = Pool#{Endpoint => ConnectionInfo},
{ok, Conn, Stream, UpdatedPool};
{error, Reason, Pool2} ->
{error, Reason, Pool2}
end;
#{connection := Conn, stream := Stream} = ConnectionInfo ->
%% Existing connection found, reuse it
%% Note: Stream is still active (quicer creates with active=true by default)
UpdatedConnectionInfo = ConnectionInfo#{last_used => erlang:system_time(second)},
UpdatedPool = Pool#{Endpoint => UpdatedConnectionInfo},
{ok, Conn, Stream, UpdatedPool}
end.
%% @doc Create a new connection to an endpoint.
-spec create_connection(binary(), binary(), binary(), map()) ->
{ok, pid(), pid(), map()} | {error, term(), map()}.
create_connection(Endpoint, NodeId, RealmId, Pool) when is_map(Pool) ->
%% Parse endpoint URL
case macula_utils:parse_url(Endpoint) of
{Host, Port} ->
?LOG_INFO("Creating connection to endpoint: ~s:~p", [Host, Port]),
%% Connect via QUIC with TLS configuration from macula_tls (v0.11.0+)
%% Use hostname-aware TLS options for certificate verification
TlsOpts = macula_tls:quic_client_opts_with_hostname(Host),
QuicOpts = merge_quic_opts([
{alpn, ["macula"]},
{idle_timeout_ms, 60000},
{keep_alive_interval_ms, 20000},
{handshake_idle_timeout_ms, 30000}
], TlsOpts),
ConnectResult = try
macula_quic:connect(Host, Port, QuicOpts, ?CONNECTION_TIMEOUT_MS)
catch
_:Error ->
{error, Error}
end,
case ConnectResult of
{ok, Conn} ->
%% Open bidirectional stream
case macula_quic:open_stream(Conn) of
{ok, Stream} ->
%% Note: quicer creates streams with active=true by default
%% Calling process (this gen_server) is already the owner
%% No need for controlling_process or explicit setopt
%% Send CONNECT message
%% Include advertise endpoint for peer-to-peer connections
LocalEndpoint = case application:get_env(macula, advertise_endpoint) of
{ok, Ep} when is_binary(Ep) -> Ep;
_ ->
%% Construct from NODE_HOST env var
NodeHost = list_to_binary(os:getenv("NODE_HOST", "localhost")),
<<"https://", NodeHost/binary, ":9443">>
end,
ConnectMsg = #{
version => <<"1.0">>,
node_id => NodeId,
realm_id => RealmId,
capabilities => [rpc], % Only RPC for endpoint connections
endpoint => LocalEndpoint
},
case send_connect_message(Stream, ConnectMsg) of
ok ->
?LOG_INFO("Connected to endpoint: ~s:~p", [Host, Port]),
{ok, Conn, Stream, Pool};
{error, Reason} ->
macula_quic:close(Stream),
macula_quic:close(Conn),
?LOG_ERROR("Handshake failed with endpoint ~s:~p: ~p",
[Host, Port, Reason]),
{error, {handshake_failed, Reason}, Pool}
end;
{error, Reason} ->
macula_quic:close(Conn),
?LOG_ERROR("Stream open failed with endpoint ~s:~p: ~p",
[Host, Port, Reason]),
{error, {stream_open_failed, Reason}, Pool}
end;
{error, Reason} ->
?LOG_ERROR("Connection failed to endpoint ~s:~p: ~p",
[Host, Port, Reason]),
{error, {connection_failed, Reason}, Pool};
{error, Type, Details} ->
?LOG_ERROR("Connection failed to endpoint ~s:~p: ~p ~p",
[Host, Port, Type, Details]),
{error, {connection_failed, {Type, Details}}, Pool};
Other ->
?LOG_ERROR("Connection failed to endpoint ~s:~p: ~p",
[Host, Port, Other]),
{error, {connection_failed, Other}, Pool}
end
end.
%% @doc Close all connections in the pool.
-spec close_all_connections(map()) -> ok.
close_all_connections(Pool) when is_map(Pool) ->
maps:foreach(
fun(_Endpoint, #{connection := Conn, stream := Stream}) ->
catch macula_quic:close(Stream),
catch macula_quic:close(Conn)
end,
Pool
),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Merge QUIC options, with second list taking precedence.
%% @private
-spec merge_quic_opts(list(), list()) -> list().
merge_quic_opts(BaseOpts, OverrideOpts) ->
lists:foldl(
fun({Key, Value}, Acc) ->
lists:keystore(Key, 1, Acc, {Key, Value})
end,
BaseOpts,
OverrideOpts
).
%% @doc Send CONNECT message through a stream.
-spec send_connect_message(pid(), map()) -> ok | {error, term()}.
send_connect_message(Stream, ConnectMsg) ->
encode_and_send(catch macula_protocol_encoder:encode(connect, ConnectMsg), Stream).
%% @private Handle encoding result and send
encode_and_send({'EXIT', Error}, _Stream) ->
{error, Error};
encode_and_send(Binary, Stream) when is_binary(Binary) ->
macula_quic:send(Stream, Binary);
encode_and_send(Other, _Stream) ->
{error, {invalid_encoding, Other}}.