%%%------------------------------------------------------------------- %%% @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}}.