Packages
macula
0.10.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_peer_system/macula_peer_connector.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Peer Connector - Establishes direct QUIC connections to remote peers (v0.8.0+).
%%%
%%% This module enables peer-to-peer communication by establishing outbound
%%% QUIC connections to arbitrary peers. Used by DHT for STORE/FIND_VALUE
%%% message propagation and by RPC/PubSub for direct delivery.
%%%
%%% == Overview ==
%%%
%%% Pattern: Connection-pooled utility module
%%% - Uses `macula_peer_connection_pool' for connection reuse
%%% - Falls back to direct connection if pool unavailable
%%% - Fire-and-forget message sending
%%%
%%% == Usage ==
%%%
%%% Used internally by:
%%% - `macula_pubsub_dht': Direct pub/sub delivery to discovered subscribers
%%% - `macula_service_registry': DHT STORE propagation to k=20 nodes
%%% - Future: Multi-hop RPC routing
%%%
%%% ```
%%% %% Send a DHT STORE message to a peer
%%% Endpoint = <<"192.168.1.100:9443">>,
%%% Message = #{
%%% key => <<"service.calculator.add">>,
%%% value => <<"192.168.1.50:9443">>,
%%% ttl => 300
%%% },
%%% ok = macula_peer_connector:send_message(Endpoint, dht_store, Message).
%%% '''
%%%
%%% == Performance Characteristics ==
%%%
%%% v0.8.0: Fire-and-forget pattern (now legacy fallback)
%%% - Creates new connection per message
%%% - Simple but inefficient for high-frequency messaging
%%%
%%% v0.10.0: Connection pooling (current)
%%% - Reuses existing connections via macula_peer_connection_pool
%%% - 1.5-2x latency improvement for repeated messaging
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_peer_connector).
-include_lib("kernel/include/logger.hrl").
-include_lib("quicer/include/quicer.hrl").
%% API
-export([
send_message/3,
send_message/4
]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Send a message to a remote peer (fire-and-forget).
%% Uses connection pool for efficiency, falls back to direct connection.
-spec send_message(binary(), atom(), map()) -> ok | {error, term()}.
send_message(Endpoint, MessageType, Message) ->
send_message(Endpoint, MessageType, Message, 5000).
%% @doc Send a message to a remote peer with custom timeout.
-spec send_message(binary(), atom(), map(), timeout()) -> ok | {error, term()}.
send_message(Endpoint, MessageType, Message, _Timeout) ->
%% Encode message
MessageBinary = macula_protocol_encoder:encode(MessageType, Message),
%% Try to use connection pool first
PoolPid = whereis(macula_peer_connection_pool),
send_via_connection(PoolPid, Endpoint, MessageBinary).
%% @private Pool not running - fall back to direct connection
send_via_connection(undefined, Endpoint, MessageBinary) ->
send_via_direct_connection(Endpoint, MessageBinary);
%% @private Pool available - use pooled connection
send_via_connection(_Pid, Endpoint, MessageBinary) ->
send_via_pool(Endpoint, MessageBinary).
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @private
%% @doc Send message using connection pool (preferred, 1.5-2x faster).
send_via_pool(Endpoint, MessageBinary) ->
ConnResult = macula_peer_connection_pool:get_connection(Endpoint),
do_pool_send(ConnResult, Endpoint, MessageBinary).
%% @private Pool connection acquired - attempt send
do_pool_send({ok, Conn, Stream}, Endpoint, MessageBinary) ->
SendResult = macula_quic:send(Stream, MessageBinary),
handle_pool_send_result(SendResult, Conn, Stream, Endpoint, MessageBinary);
%% @private Pool connection failed - fall back to direct
do_pool_send({error, Reason}, Endpoint, MessageBinary) ->
?LOG_DEBUG("Pool connection failed: ~p, using direct", [Reason]),
send_via_direct_connection(Endpoint, MessageBinary).
%% @private Send succeeded - return connection to pool
handle_pool_send_result(ok, Conn, Stream, Endpoint, _MessageBinary) ->
macula_peer_connection_pool:return_connection(Endpoint, {Conn, Stream}),
ok;
%% @private Send failed - invalidate and retry with direct
handle_pool_send_result({error, Reason}, _Conn, _Stream, Endpoint, MessageBinary) ->
macula_peer_connection_pool:invalidate(Endpoint),
?LOG_WARNING("Pool send failed: ~p, falling back to direct", [Reason]),
send_via_direct_connection(Endpoint, MessageBinary).
%% @private
%% @doc Send message via direct connection (fallback, creates new connection).
send_via_direct_connection(Endpoint, MessageBinary) ->
ParseResult = parse_endpoint(Endpoint),
do_direct_send(ParseResult, Endpoint, MessageBinary).
%% @private Endpoint parsed successfully
do_direct_send({ok, Host, Port}, _Endpoint, MessageBinary) ->
send_via_quic(Host, Port, MessageBinary, 5000);
%% @private Invalid endpoint format
do_direct_send({error, Reason}, Endpoint, _MessageBinary) ->
?LOG_ERROR("Invalid endpoint ~p: ~p", [Endpoint, Reason]),
{error, {invalid_endpoint, Reason}}.
%% @private
%% @doc Parse endpoint string into host and port.
parse_endpoint(Endpoint) when is_binary(Endpoint) ->
parse_endpoint(binary_to_list(Endpoint));
parse_endpoint(Endpoint) when is_list(Endpoint) ->
SplitResult = string:split(Endpoint, ":"),
do_parse_endpoint(SplitResult).
%% @private Valid host:port format - parse port
do_parse_endpoint([Host, PortStr]) ->
parse_port(Host, PortStr);
%% @private Invalid format
do_parse_endpoint(_) ->
{error, invalid_format}.
%% @private Parse port string to integer
parse_port(Host, PortStr) ->
parse_port_result(Host, catch list_to_integer(PortStr)).
%% @private Port parsed successfully
parse_port_result(Host, Port) when is_integer(Port), Port > 0, Port < 65536 ->
{ok, Host, Port};
%% @private Invalid port value or parse error
parse_port_result(_Host, _) ->
{error, invalid_port}.
%% @private
%% @doc Send message via direct QUIC connection (legacy fallback).
send_via_quic(Host, Port, MessageBinary, Timeout) ->
ConnectOpts = [
{alpn, ["macula"]},
{verify, none},
{idle_timeout_ms, 60000},
{keep_alive_interval_ms, 20000},
{handshake_idle_timeout_ms, 30000}
],
ConnResult = macula_quic:connect(Host, Port, ConnectOpts, Timeout),
do_quic_connect(ConnResult, MessageBinary).
%% @private Connection established - open stream
do_quic_connect({ok, Conn}, MessageBinary) ->
StreamResult = macula_quic:open_stream(Conn),
do_quic_stream(StreamResult, Conn, MessageBinary);
%% @private Transport down
do_quic_connect({error, transport_down, _Details}, _MessageBinary) ->
{error, {connect_failed, transport_down}};
%% @private Connection failed
do_quic_connect({error, Reason}, _MessageBinary) ->
{error, {connect_failed, Reason}}.
%% @private Stream opened - send message
do_quic_stream({ok, Stream}, Conn, MessageBinary) ->
SendResult = macula_quic:send(Stream, MessageBinary),
do_quic_send(SendResult, Conn, Stream);
%% @private Stream open failed
do_quic_stream({error, Reason}, Conn, _MessageBinary) ->
quicer:async_shutdown_connection(Conn, ?QUIC_CONNECTION_SHUTDOWN_FLAG_NONE, 0),
{error, {stream_failed, Reason}}.
%% @private Send succeeded - graceful shutdown
do_quic_send(ok, Conn, Stream) ->
timer:sleep(50),
quicer:async_shutdown_stream(Stream, ?QUIC_STREAM_SHUTDOWN_FLAG_GRACEFUL, 0),
quicer:async_shutdown_connection(Conn, ?QUIC_CONNECTION_SHUTDOWN_FLAG_NONE, 0),
ok;
%% @private Send failed - abort stream
do_quic_send({error, Reason}, Conn, Stream) ->
quicer:async_shutdown_stream(Stream, ?QUIC_STREAM_SHUTDOWN_FLAG_ABORT, 0),
quicer:async_shutdown_connection(Conn, ?QUIC_CONNECTION_SHUTDOWN_FLAG_NONE, 0),
{error, {send_failed, Reason}}.