Packages

macula

0.48.6
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_connection.erl
Raw

src/macula_connection.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Macula Connection - QUIC Transport Layer (v0.7.0+).
%%%
%%% This module manages the low-level QUIC connection lifecycle and
%%% message transport for mesh participants.
%%%
%%% Responsibilities:
%%% - Establish and maintain QUIC connection
%%% - Send messages via QUIC stream
%%% - Receive and route incoming messages to handlers
%%% - Handle connection errors and reconnection
%%% - Message encoding/decoding and buffering
%%%
%%% Renamed from macula_connection in v0.7.0 for clarity:
%%% - macula_connection = QUIC transport (this module - low-level)
%%% - macula_peer = mesh participant (high-level API)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_connection).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-include("macula_config.hrl").
-include("macula_connection.hrl").
%% API
-export([start_link/2, send_message/3, send_message_async/3, get_status/1, default_config/0]).
-export([hostname_from_node/0]).
-export([build_node_identity/1]).
%% Bridge system API (v0.13.0+)
%% Used by macula_bridge_node for parent mesh connections.
%% Uses bridge_rpc and bridge_data message types for communication.
-export([close/1, call/4, send/2, connect/3]).
%% API for testing
-export([decode_messages/2]).
%% Internal exports
-export([start_keepalive_timer/1]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
%% Test exports for bootstrap lookup helpers moved to macula_connection_dispatch
-define(SERVER, ?MODULE).
-define(INITIAL_RETRY_DELAY_MS, 2000).
-define(MAX_RETRY_DELAY_MS, 120000).
%%%===================================================================
%%% API
%%%===================================================================
-spec start_link(binary(), map()) -> {ok, pid()} | {error, term()}.
start_link(Url, Opts) ->
gen_server:start_link(?MODULE, {Url, Opts}, []).
-spec default_config() -> map().
default_config() ->
#{
keepalive_enabled => true,
keepalive_interval => 30000 %% 30 seconds
}.
-spec send_message(pid(), atom(), map()) -> ok | {error, term()}.
send_message(Pid, Type, Msg) ->
gen_server:call(Pid, {send_message, Type, Msg}, 5000).
%% @doc Send message asynchronously (fire-and-forget).
%% Use for operations where blocking is unacceptable and failures can be tolerated.
%% The message will be sent if connected, silently dropped if not.
-spec send_message_async(pid(), atom(), map()) -> ok.
send_message_async(Pid, Type, Msg) ->
gen_server:cast(Pid, {send_message_async, Type, Msg}).
-spec get_status(pid()) -> connecting | connected | disconnected | error.
get_status(Pid) ->
gen_server:call(Pid, get_status, 5000).
%%%===================================================================
%%% Bridge System API (v0.13.0+)
%%% These functions provide a simplified API for bridge-to-bridge
%%% connections, using the underlying QUIC transport layer.
%%%===================================================================
%% @doc Close a connection gracefully.
%% Stops the gen_server which triggers proper QUIC cleanup in terminate/2.
%% Uses catch expression to handle race conditions where process may already be dead.
-spec close(pid()) -> ok.
close(Pid) when is_pid(Pid) ->
_ = (catch gen_server:stop(Pid, normal, 5000)),
ok;
close(_) ->
ok.
%% @doc Make an RPC-style call over a connection.
%% Sends a call message and waits for a reply. Uses the underlying
%% send_message API with type 'bridge_rpc'.
-spec call(pid(), binary(), map(), map()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args, Opts) when is_pid(Pid) ->
Timeout = maps:get(timeout, Opts, 5000),
CallId = crypto:strong_rand_bytes(16),
Msg = #{
procedure => Procedure,
args => Args,
call_id => CallId,
timeout => Timeout
},
%% Use synchronous send which validates connection state
case send_message(Pid, bridge_rpc, Msg) of
ok -> {ok, sent}; % Bridge RPC is fire-and-forget at transport level
{error, Reason} -> {error, Reason}
end;
call(_, _, _, _) ->
{error, invalid_connection}.
%% @doc Send a message over a connection asynchronously.
%% Fire-and-forget delivery - returns immediately.
-spec send(pid(), term()) -> ok | {error, term()}.
send(Pid, Message) when is_pid(Pid), is_map(Message) ->
send_message_async(Pid, bridge_data, Message),
ok;
send(Pid, Message) when is_pid(Pid) ->
%% Wrap non-map messages
send_message_async(Pid, bridge_data, #{payload => Message}),
ok;
send(_, _) ->
{error, invalid_connection}.
%% @doc Connect to a remote endpoint.
%% Creates a new QUIC connection to the specified host:port.
-spec connect(binary(), pos_integer(), map()) -> {ok, pid()} | {error, term()}.
connect(Host, Port, Opts) when is_binary(Host), is_integer(Port), Port > 0 ->
Url = iolist_to_binary([<<"quic://">>, Host, <<":">>, integer_to_binary(Port)]),
start_link(Url, Opts);
connect(Host, Port, Opts) when is_list(Host) ->
connect(list_to_binary(Host), Port, Opts);
connect(_, _, _) ->
{error, invalid_arguments}.
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init({Url, Opts}) ->
?LOG_INFO("[Connection] Starting for ~s", [Url]),
%% Parse URL to extract host and port
{Host, Port} = parse_url(Url),
%% Get realm (required)
Realm = get_realm_from_opts(Opts),
%% Generate or get node ID
NodeId = maps:get(node_id, Opts, generate_node_id()),
%% Get peer_id from opts (set by macula_peer_system)
PeerId = maps:get(peer_id, Opts, erlang:unique_integer([monotonic, positive])),
State = #state{
url = Url,
opts = Opts#{host => Host, port => Port},
node_id = NodeId,
realm = Realm,
peer_id = PeerId,
connection = undefined,
stream = undefined,
status = connecting,
recv_buffer = <<>>,
keepalive_timer = undefined
},
%% Initiate connection asynchronously
self() ! connect,
{ok, State}.
handle_call({send_message, Type, Msg}, _From, #state{status = connected, stream = Stream} = State) ->
case macula_connection_dispatch:send_message_raw(Type, Msg, Stream) of
ok ->
{reply, ok, State};
{error, Reason} ->
{reply, {error, Reason}, State}
end;
handle_call({send_message, _Type, _Msg}, _From, State) ->
{reply, {error, not_connected}, State};
handle_call(get_status, _From, State) ->
{reply, State#state.status, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
%% @doc Handle async send message - fire and forget.
%% Sends if connected, logs warning and drops if not.
handle_cast({send_message_async, Type, Msg}, #state{status = connected, stream = Stream} = State) ->
case macula_connection_dispatch:send_message_raw(Type, Msg, Stream) of
ok ->
{noreply, State};
{error, Reason} ->
?LOG_WARNING("Async send failed for type ~p: ~p", [Type, Reason]),
{noreply, State}
end;
handle_cast({send_message_async, Type, _Msg}, State) ->
?LOG_DEBUG("Async send dropped (not connected): type=~p, status=~p", [Type, State#state.status]),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%% Already connecting — ignore duplicate connect messages
handle_info(connect, #state{connect_in_flight = true} = State) ->
{noreply, State};
%% Already connected — ignore
handle_info(connect, #state{status = connected} = State) ->
{noreply, State};
%% Not connected, not in flight — attempt connection
handle_info(connect, State) ->
?LOG_DEBUG("[Connection] Attempting async connect to ~s", [State#state.url]),
spawn_connect(State),
{noreply, State#state{connect_in_flight = true}};
handle_info({connect_result, {ok, Conn, Stream}}, State) ->
%% Ownership was already transferred by the spawned process before it exited.
%% Now set active mode so we receive {quic, ...} messages.
macula_quic:setopt(Stream, active, true),
#{host := Host, port := Port} = State#state.opts,
case complete_connection_setup(Conn, Stream, Host, Port, State) of
{ok, State2} ->
?LOG_INFO("[Connection] Successfully connected to ~s", [State#state.url]),
{noreply, State2#state{connect_in_flight = false,
retry_delay_ms = ?INITIAL_RETRY_DELAY_MS}};
{error, Reason} ->
NextDelay = State#state.retry_delay_ms,
?LOG_ERROR("[Connection] Connection setup failed for ~s: ~p, retrying in ~.1fs",
[State#state.url, Reason, NextDelay / 1000]),
erlang:send_after(NextDelay, self(), connect),
{noreply, State#state{status = error, connect_in_flight = false,
retry_delay_ms = min(NextDelay * 2, ?MAX_RETRY_DELAY_MS)}}
end;
handle_info({connect_result, {error, Reason}}, State) ->
NextDelay = State#state.retry_delay_ms,
?LOG_ERROR("[Connection] Connection setup failed for ~s: ~p, retrying in ~.1fs",
[State#state.url, Reason, NextDelay / 1000]),
erlang:send_after(NextDelay, self(), connect),
{noreply, State#state{status = error, connect_in_flight = false,
retry_delay_ms = min(NextDelay * 2, ?MAX_RETRY_DELAY_MS)}};
handle_info({quic, Data, Stream, _Props}, State) when is_binary(Data) ->
%% Received data from QUIC stream
MainStream = State#state.stream,
?LOG_DEBUG("[Connection] QUIC data received: ~p bytes, from_stream=~p, main_stream=~p, match=~p",
[byte_size(Data), Stream, MainStream, Stream =:= MainStream]),
handle_stream_data(Stream =:= MainStream, Data, Stream, State);
%% Handle QUIC control messages on the MAIN stream — trigger reconnect
handle_info({quic, ControlMsg, Stream, _Props}, #state{stream = Stream} = State) when is_atom(ControlMsg) ->
handle_quic_control_message(ControlMsg, Stream, State);
%% Handle QUIC control messages on OTHER streams (temp DHT streams) — ignore
handle_info({quic, ControlMsg, OtherStream, _Props}, State) when is_atom(ControlMsg) ->
?LOG_DEBUG("[Connection] Ignoring control message ~p on non-main stream ~p", [ControlMsg, OtherStream]),
{noreply, State};
%% Handle keep-alive tick - send PING message
handle_info(keepalive_tick, #state{status = connected, stream = Stream} = State) ->
%% Send PING message
PingMsg = #{timestamp => erlang:system_time(millisecond)},
case macula_connection_dispatch:send_message_raw(ping, PingMsg, Stream) of
ok ->
?LOG_DEBUG("Keep-alive PING sent"),
%% Restart timer for next keep-alive
StateWithTimer = start_keepalive_timer(State),
{noreply, StateWithTimer};
{error, Reason} ->
?LOG_WARNING("Failed to send keep-alive PING: ~p", [Reason]),
%% Connection might be dead - let it retry
{noreply, State}
end;
%% Ignore keep-alive tick if not connected
handle_info(keepalive_tick, State) ->
{noreply, State};
%% Mesh peer lifecycle events (from gateway pub/sub via _mesh.peer.* topics)
handle_info({mesh_peer_connected, PeerInfo}, State) ->
PeerNodeId = maps:get(<<"node_id">>, PeerInfo, maps:get(node_id, PeerInfo, <<>>)),
?LOG_INFO("[Connection] Peer connected: ~s", [PeerNodeId]),
macula_connection_dispatch:add_discovered_peers([PeerInfo], State#state.node_id),
macula_connection_dispatch:notify_mesh_lifecycle_observers(mesh_peer_connected, PeerInfo),
{noreply, State};
handle_info({mesh_peer_disconnected, PeerInfo}, State) ->
NodeId = maps:get(<<"node_id">>, PeerInfo, maps:get(node_id, PeerInfo, undefined)),
RawNodeId = macula_connection_dispatch:normalize_peer_node_id(NodeId),
?LOG_INFO("[Connection] Peer disconnected: ~s",
[binary:encode_hex(RawNodeId)]),
macula_connection_dispatch:remove_peer_from_routing_table(RawNodeId),
macula_connection_dispatch:notify_mesh_lifecycle_observers(mesh_peer_disconnected, PeerInfo),
{noreply, State};
handle_info(_Info, State) ->
%% Ignore unexpected messages
?LOG_DEBUG("Unhandled handle_info message: ~p", [_Info]),
{noreply, State}.
terminate(_Reason, #state{stream = Stream, connection = Conn}) ->
?LOG_INFO("Connection manager terminating"),
%% Do NOT invalidate pool — pool detects dead connections via is_connection_alive.
%% Clean up QUIC resources
catch macula_quic:close(Stream),
catch macula_quic:close(Conn),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Dispatch stream data based on validity (pattern matching on boolean).
handle_stream_data(true, Data, _Stream, State) ->
?LOG_DEBUG("[Connection] handle_stream_data(true) ENTRY: ~p bytes", [byte_size(Data)]),
handle_received_data(Data, State);
handle_stream_data(false, _Data, Stream, State) ->
?LOG_WARNING("Received data from unknown stream: ~p", [Stream]),
{noreply, State}.
%% @doc Handle QUIC control messages - trigger reconnect on connection/stream closure.
-spec handle_quic_control_message(atom(), pid(), #state{}) -> {noreply, #state{}}.
%% Stream closed - trigger reconnection
handle_quic_control_message(peer_send_shutdown, _Stream, State) ->
?LOG_WARNING("QUIC stream closed by peer (peer_send_shutdown), reconnecting..."),
trigger_reconnect(State);
handle_quic_control_message(peer_send_aborted, _Stream, State) ->
?LOG_WARNING("QUIC stream aborted by peer, reconnecting..."),
trigger_reconnect(State);
handle_quic_control_message(send_shutdown_complete, _Stream, State) ->
?LOG_WARNING("QUIC send shutdown complete, reconnecting..."),
trigger_reconnect(State);
handle_quic_control_message(shutdown, _Stream, State) ->
?LOG_WARNING("QUIC shutdown, reconnecting..."),
trigger_reconnect(State);
handle_quic_control_message(closed, _Stream, State) ->
?LOG_WARNING("QUIC connection closed, reconnecting..."),
trigger_reconnect(State);
%% Other control messages - log and ignore
handle_quic_control_message(ControlMsg, _Stream, State) ->
?LOG_DEBUG("Ignoring QUIC control message: ~p", [ControlMsg]),
{noreply, State}.
%% @doc Trigger reconnection by cleaning up and scheduling reconnect.
-spec trigger_reconnect(#state{}) -> {noreply, #state{}}.
trigger_reconnect(#state{stream = Stream, connection = Conn, keepalive_timer = Timer} = State) ->
%% Cancel keep-alive timer
case Timer of
undefined -> ok;
_ -> erlang:cancel_timer(Timer)
end,
%% Do NOT invalidate pool here — another macula_connection instance may
%% have seeded the pool with ITS connection. The pool's is_connection_alive
%% check will detect dead connections on the next get_connection call.
%% Clean up old connection resources
catch macula_quic:close(Stream),
catch macula_quic:close(Conn),
%% Schedule reconnect with backoff (reset to initial since this was a working connection)
erlang:send_after(?INITIAL_RETRY_DELAY_MS, self(), connect),
%% Update state to disconnected
NewState = State#state{
connection = undefined,
stream = undefined,
status = disconnected,
recv_buffer = <<>>,
keepalive_timer = undefined,
connect_in_flight = false,
retry_delay_ms = ?INITIAL_RETRY_DELAY_MS
},
{noreply, NewState}.
%% @doc Spawn async QUIC connection attempt so the gen_server remains responsive
%% to calls (e.g., get_status) while the QUIC handshake is in progress.
-spec spawn_connect(#state{}) -> pid().
spawn_connect(State) ->
Self = self(),
#{host := Host, port := Port} = State#state.opts,
QuicOpts = build_quic_opts(Host),
spawn(fun() ->
Result = try
case attempt_quic_connection(Host, Port, QuicOpts) of
{ok, Conn, Stream} ->
%% Transfer ownership to gen_server BEFORE this process exits.
%% If we don't, quicer closes the stream/connection when this
%% spawned process terminates (it's the controlling process).
ok = macula_quic:controlling_process(Conn, Self),
ok = macula_quic:controlling_process(Stream, Self),
{ok, Conn, Stream};
Other ->
Other
end
catch
_:Reason -> {error, {crashed, Reason}}
end,
Self ! {connect_result, Result}
end).
%% @doc Build QUIC connection options with TLS configuration.
%% Uses macula_tls module for centralized TLS settings (v0.11.0+).
%% Hostname is used for TLS hostname verification in production mode.
-spec build_quic_opts(Host :: string() | binary()) -> list().
build_quic_opts(Host) ->
%% Get TLS options with hostname verification from centralized module
TlsOpts = macula_tls:quic_client_opts_with_hostname(Host),
%% Merge with QUIC-specific options
BaseOpts = [
{alpn, ["macula"]},
{idle_timeout_ms, 60000},
{keep_alive_interval_ms, 20000},
{handshake_idle_timeout_ms, 30000}
],
merge_opts(BaseOpts, TlsOpts).
%% @doc Merge two option lists, with second list taking precedence.
-spec merge_opts(list(), list()) -> list().
merge_opts(BaseOpts, OverrideOpts) ->
lists:foldl(
fun({Key, Value}, Acc) ->
lists:keystore(Key, 1, Acc, {Key, Value})
end,
BaseOpts,
OverrideOpts
).
%% @doc Attempt QUIC connection and stream setup
-spec attempt_quic_connection(string(), integer(), list()) ->
{ok, pid(), pid()} | {error, term()}.
attempt_quic_connection(Host, Port, QuicOpts) ->
case safe_quic_connect(Host, Port, QuicOpts) of
{ok, Conn} ->
setup_bidirectional_stream(Conn);
{error, _Reason} = Error ->
Error;
{error, _Type, _Details} = Error ->
{error, {connection_failed, Error}};
Other ->
{error, {connection_failed, Other}}
end.
%% @doc Safe QUIC connect with error handling.
%% NIF boundary - quicer can throw exceptions, convert to tagged tuples.
-spec safe_quic_connect(string(), integer(), list()) ->
{ok, pid()} | {error, term()}.
safe_quic_connect(Host, Port, QuicOpts) ->
handle_quic_connect_result(catch macula_quic:connect(Host, Port, QuicOpts, ?CONNECTION_TIMEOUT_MS)).
%% @private Handle QUIC connect result, converting exceptions to error tuples.
handle_quic_connect_result({'EXIT', Error}) ->
{error, {connection_failed, Error}};
handle_quic_connect_result({ok, _Conn} = Result) ->
Result;
handle_quic_connect_result({error, _Reason} = Result) ->
Result;
handle_quic_connect_result({error, Type, Details}) ->
{error, {quic_transport_error, Type, Details}};
handle_quic_connect_result(Other) ->
{error, {unexpected_connect_result, Other}}.
%% @doc Open and configure bidirectional stream
-spec setup_bidirectional_stream(reference()) ->
{ok, reference(), reference()} | {error, term()}.
setup_bidirectional_stream(Conn) ->
case macula_quic:open_stream(Conn) of
{ok, Stream} ->
configure_stream_active_mode(Conn, Stream);
{error, Reason} ->
macula_quic:close(Conn),
{error, {stream_open_failed, Reason}}
end.
%% @doc Verify stream is usable after opening.
%% NOTE: active mode is NOT set here because this runs in a spawned process.
%% The gen_server takes ownership and sets active mode in handle_info({connect_result, ...}).
%% Setting active here would make the spawned process the stream owner, and when it
%% terminates, quicer closes the stream — causing {error, closed} on subsequent sends.
-spec configure_stream_active_mode(reference(), reference()) ->
{ok, reference(), reference()} | {error, term()}.
configure_stream_active_mode(Conn, Stream) ->
{ok, Conn, Stream}.
%% @doc Complete connection setup with handshake and DHT registration
-spec complete_connection_setup(reference(), reference(), string(), integer(), #state{}) ->
{ok, #state{}} | {error, term()}.
complete_connection_setup(Conn, Stream, Host, Port, State) ->
ConnectMsg = build_connect_message(State),
case macula_connection_dispatch:send_message_raw(connect, ConnectMsg, Stream) of
ok ->
?LOG_INFO("Connected to Macula mesh: ~s:~p", [Host, Port]),
%% Add server to DHT routing table immediately so STORE/FIND_VALUE
%% operations have a target. In client mode, this is the ONLY peer
%% in the routing table — without it, DHT operations go nowhere.
add_server_to_routing_table(Host, Port),
%% Seed connection pool so DHT reuses this QUIC connection
seed_connection_pool(Host, Port, Conn),
%% Start keep-alive timer if enabled
ConnectedState = State#state{
connection = Conn,
stream = Stream,
status = connected
},
StateWithKeepalive = start_keepalive_timer(ConnectedState),
%% Replay subscriptions so the gateway knows about them on the new stream.
%% After reconnect, the gateway lost all stream subscriptions from the old
%% connection. The pubsub_handler still has the application's subscriptions.
replay_peer_subscriptions(StateWithKeepalive),
{ok, StateWithKeepalive};
{error, Reason} ->
macula_quic:close(Stream),
macula_quic:close(Conn),
{error, {handshake_failed, Reason}}
end.
%% @doc Build CONNECT protocol message
-spec build_connect_message(#state{}) -> map().
build_connect_message(State) ->
LocalEndpoint = get_advertise_endpoint(),
Base = #{
version => <<"1.0">>,
node_id => State#state.node_id,
realm_id => State#state.realm,
capabilities => [pubsub, rpc],
endpoint => LocalEndpoint
},
case build_node_identity(State#state.opts) of
Identity when map_size(Identity) > 0 -> Base#{identity => Identity};
_ -> Base
end.
%% @doc Build node identity from opts or environment variables.
%% Opts key `identity' takes precedence. Falls back to HECATE_GEO_*
%% and MACULA_GEO_* env vars (same pattern as relay identity).
-spec build_node_identity(map()) -> map().
build_node_identity(Opts) ->
case maps:get(identity, Opts, undefined) of
Identity when is_map(Identity), map_size(Identity) > 0 ->
Identity;
_ ->
collect_identity_from_env()
end.
collect_identity_from_env() ->
Pairs = [
{city, [{"HECATE_GEO_CITY", bin}, {"MACULA_GEO_CITY", bin}]},
{country, [{"HECATE_GEO_COUNTRY", bin}, {"MACULA_GEO_COUNTRY", bin}]},
{lat, [{"HECATE_GEO_LAT", float}, {"MACULA_GEO_LAT", float}]},
{lng, [{"HECATE_GEO_LNG", float}, {"MACULA_GEO_LNG", float}]},
{owner, [{"HECATE_OWNER_NAME", bin}]},
{site, [{"HECATE_SITE_NAME", bin}]}
],
maps:from_list([{K, V} || {K, Sources} <- Pairs,
V <- [env_first(Sources)],
V =/= null]).
env_first([]) -> null;
env_first([{Key, Type} | Rest]) ->
case os:getenv(Key) of
false -> env_first(Rest);
"" -> env_first(Rest);
Val -> convert_env(Val, Type)
end.
convert_env(Val, bin) -> list_to_binary(Val);
convert_env(Val, float) ->
try list_to_float(Val)
catch error:badarg ->
try float(list_to_integer(Val))
catch error:badarg -> null
end
end.
%% @doc Get endpoint to advertise for mesh connections
-spec get_advertise_endpoint() -> binary().
get_advertise_endpoint() ->
case application:get_env(macula, advertise_endpoint) of
{ok, Endpoint} when is_binary(Endpoint) ->
Endpoint;
_ ->
construct_default_endpoint()
end.
%% @doc Notify the peer's pubsub_handler to re-send SUBSCRIBE messages after reconnect.
-spec replay_peer_subscriptions(#state{}) -> ok.
replay_peer_subscriptions(#state{realm = Realm, peer_id = PeerId}) ->
case gproc:lookup_local_name({pubsub_handler, Realm, PeerId}) of
undefined ->
?LOG_DEBUG("[Connection] No pubsub_handler for replay (realm=~s, peer_id=~p)",
[Realm, PeerId]);
PubSubPid ->
gen_server:cast(PubSubPid, replay_subscriptions)
end,
ok.
%% @doc Add the connected server to the local DHT routing table.
%% In client mode, this is the ONLY peer — without it, STORE and FIND_VALUE
%% have no target and DHT operations silently fail (empty results).
-spec add_server_to_routing_table(string(), integer()) -> ok.
add_server_to_routing_table(Host, Port) ->
Endpoint = iolist_to_binary([<<"https://">>, list_to_binary(Host), <<":">>, integer_to_binary(Port)]),
ServerNodeId = crypto:hash(sha256, Endpoint),
ServerNodeInfo = #{
node_id => ServerNodeId,
address => Endpoint,
endpoint => Endpoint
},
case whereis(macula_routing_server) of
undefined -> ok;
_Pid ->
case macula_routing_server:add_node(macula_routing_server, ServerNodeInfo) of
ok ->
?LOG_INFO("[Connection] Added server to routing table: ~s", [Endpoint]);
{error, Reason} ->
?LOG_WARNING("[Connection] Failed to add server to routing table: ~p", [Reason])
end
end,
ok.
%% @doc Construct default endpoint from environment and application config.
%% Port resolution order: MACULA_QUIC_PORT env var, application:get_env(macula, quic_port), 9443.
-spec construct_default_endpoint() -> binary().
construct_default_endpoint() ->
NodeHost = get_hostname_from_env(),
Port = get_quic_port_binary(),
<<"https://", NodeHost/binary, ":", Port/binary>>.
%% @doc Resolve QUIC port: env var first, then app config, then default 9443.
-spec get_quic_port_binary() -> binary().
get_quic_port_binary() ->
case os:getenv("MACULA_QUIC_PORT") of
false ->
case application:get_env(macula, quic_port) of
{ok, P} when is_integer(P) -> integer_to_binary(P);
{ok, P} when is_list(P) -> list_to_binary(P);
_ -> <<"9443">>
end;
EnvPort -> list_to_binary(EnvPort)
end.
%% @doc Get hostname from environment variables.
%% Checks MACULA_HOSTNAME first (Docker), then NODE_HOST, then HOSTNAME, fallback to localhost.
-spec get_hostname_from_env() -> binary().
get_hostname_from_env() ->
get_hostname_from_env([
"MACULA_HOSTNAME", %% Docker compose sets this
"NODE_HOST", %% Legacy/alternative
"HOSTNAME" %% Standard shell variable
]).
%% @doc Try environment variables in order, then node name, then localhost.
-spec get_hostname_from_env([string()]) -> binary().
get_hostname_from_env([]) ->
hostname_from_node();
get_hostname_from_env([EnvVar | Rest]) ->
case os:getenv(EnvVar) of
false -> get_hostname_from_env(Rest);
Value -> list_to_binary(Value)
end.
%% @private Extract hostname from Erlang node name (e.g. hecate@beam00.lab -> beam00.lab).
hostname_from_node() ->
case atom_to_list(node()) of
"nonode@nohost" -> <<"localhost">>;
NodeStr ->
case string:split(NodeStr, "@") of
[_, Host] -> list_to_binary(Host);
_ -> <<"localhost">>
end
end.
%% @doc Handle received data from the stream.
-spec handle_received_data(binary(), #state{}) -> {noreply, #state{}}.
handle_received_data(Data, State) ->
%% Append to receive buffer
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
?LOG_DEBUG("[Connection] handle_received_data ENTRY: data=~p bytes, buffer=~p bytes", [byte_size(Data), byte_size(Buffer)]),
%% Try to decode messages
{Messages, RemainingBuffer} = decode_messages(Buffer, []),
?LOG_DEBUG("[Connection] decode_messages returned: ~p messages, remaining=~p bytes", [length(Messages), byte_size(RemainingBuffer)]),
%% Process each message
State2 = lists:foldl(fun macula_connection_dispatch:process_message/2, State, Messages),
{noreply, State2#state{recv_buffer = RemainingBuffer}}.
%% @doc Decode all complete messages from buffer.
-spec decode_messages(binary(), list()) -> {list(), binary()}.
decode_messages(Buffer, Acc) when byte_size(Buffer) < 8 ->
%% Not enough for header
{lists:reverse(Acc), Buffer};
decode_messages(<<_Version:8, _TypeId:8, _Flags:8, _Reserved:8,
PayloadLen:32/big-unsigned, Rest/binary>> = Buffer, Acc) ->
case byte_size(Rest) of
ActualLen when ActualLen >= PayloadLen ->
%% We have a complete message
?LOG_DEBUG("[Connection] DECODING message: buffer=~p bytes, payload_len=~p", [byte_size(Buffer), PayloadLen]),
case macula_protocol_decoder:decode(Buffer) of
{ok, {Type, Msg}} ->
%% Skip this message and continue
?LOG_DEBUG("[Connection] DECODED OK: type=~p", [Type]),
<<_:8/binary, _Payload:PayloadLen/binary, Remaining/binary>> = Buffer,
decode_messages(Remaining, [{Type, Msg} | Acc]);
{error, Reason} ->
%% Decode error, skip this message - LOG THE ERROR
<<_:8, TypeIdByte:8, _:6/binary, _Payload:PayloadLen/binary, Remaining/binary>> = Buffer,
?LOG_WARNING("[Connection] DECODE ERROR: type_id=~p (0x~.16B), reason=~p, payload_len=~p",
[TypeIdByte, TypeIdByte, Reason, PayloadLen]),
decode_messages(Remaining, Acc)
end;
_ ->
%% Incomplete message, wait for more data
{lists:reverse(Acc), Buffer}
end.
%% @doc Parse URL to extract host and port.
-spec parse_url(binary()) -> {string(), integer()}.
parse_url(Url) when is_binary(Url) ->
parse_url(binary_to_list(Url));
parse_url("https://" ++ Rest) ->
parse_host_port(Rest, 443);
parse_url("http://" ++ Rest) ->
parse_host_port(Rest, 80);
parse_url(Url) ->
error({invalid_url, Url}).
%% @doc Parse host and port from URL remainder, using default port if not specified.
-spec parse_host_port(string(), integer()) -> {string(), integer()}.
parse_host_port(HostPort, DefaultPort) ->
case string:split(HostPort, ":") of
[Host, PortStr] -> {Host, list_to_integer(PortStr)};
[Host] -> {Host, DefaultPort}
end.
%% @doc Generate a random node ID.
-spec generate_node_id() -> binary().
generate_node_id() ->
macula_utils:generate_node_id().
%% @doc Get realm from options and normalize.
-spec get_realm_from_opts(map()) -> binary().
get_realm_from_opts(Opts) ->
normalize_realm(maps:get(realm, Opts, undefined)).
%% @doc Normalize realm to binary.
-spec normalize_realm(undefined | binary() | list() | atom()) -> binary().
normalize_realm(undefined) ->
error({missing_required_option, realm});
normalize_realm(Realm) when is_binary(Realm) ->
Realm;
normalize_realm(Realm) when is_list(Realm) ->
list_to_binary(Realm);
normalize_realm(Realm) when is_atom(Realm) ->
atom_to_binary(Realm).
%% @doc Start keep-alive timer if enabled in options.
-spec start_keepalive_timer(#state{}) -> #state{}.
start_keepalive_timer(#state{opts = Opts, keepalive_timer = OldTimer} = State) ->
%% Cancel existing timer if present
case OldTimer of
undefined -> ok;
_ -> erlang:cancel_timer(OldTimer)
end,
%% Check if keep-alive is enabled
Enabled = maps:get(keepalive_enabled, Opts, true),
start_keepalive_timer(Enabled, State).
%% Keep-alive disabled
start_keepalive_timer(false, State) ->
State#state{keepalive_timer = undefined};
%% Keep-alive enabled - start timer
start_keepalive_timer(true, #state{opts = Opts} = State) ->
Interval = maps:get(keepalive_interval, Opts, 30000),
TimerRef = erlang:send_after(Interval, self(), keepalive_tick),
State#state{keepalive_timer = TimerRef}.
%%%===================================================================
%%% Connection Pool Seeding
%%%===================================================================
%% @doc Seed the peer connection pool with the main QUIC connection.
%% This allows DHT operations (STORE, FIND_VALUE) to reuse the existing
%% connection instead of opening new ones that trigger rate limiting.
-spec seed_connection_pool(string(), integer(), pid()) -> ok.
seed_connection_pool(Host, Port, Conn) ->
PoolKey = pool_key(Host, Port),
case whereis(macula_peer_connection_pool) of
undefined ->
?LOG_DEBUG("[Connection] Pool not running, skipping seed for ~s", [PoolKey]),
ok;
_Pid ->
?LOG_INFO("[Connection] Seeding connection pool: ~s", [PoolKey]),
macula_peer_connection_pool:put(PoolKey, Conn),
ok
end.
%% @doc Build pool key matching the format used by DHT/peer_connector.
%% Returns a bare "host:port" binary with no scheme prefix.
-spec pool_key(string(), integer()) -> binary().
pool_key(Host, Port) ->
iolist_to_binary([Host, ":", integer_to_list(Port)]).