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
macula src macula_gateway_system macula_gateway_quic_server.erl
Raw

src/macula_gateway_system/macula_gateway_quic_server.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% QUIC Transport Layer Gen_Server
%%%
%%% Handles all QUIC transport operations for the gateway:
%%% - Owns QUIC listener
%%% - Receives {quic, ...} events
%%% - Decodes protocol messages
%%% - Routes messages to gateway for business logic
%%%
%%% This separation follows proper OTP design:
%%% - One process, one responsibility (transport vs routing)
%%% - Clean fault isolation (QUIC crashes don't crash gateway)
%%% - Proper supervision (supervisor can restart independently)
%%% - Testability (can test transport in isolation)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_quic_server).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% Suppress warnings for functions only used in tests
-compile({nowarn_unused_function, [parse_endpoint/1, resolve_host/2]}).
%% API
-export([start_link/1, set_gateway/2]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
%% Helper functions (exported for testing)
-ifdef(TEST).
-export([parse_endpoint/1, resolve_host/2, complete_handshake/1,
accept_streams/1, register_next_connection/1]).
-endif.
-record(state, {
listener :: pid() | undefined,
gateway :: pid() | undefined,
node_id :: binary(),
port :: inet:port_number(),
realm :: binary(),
buffer = <<>> :: binary()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the QUIC server gen_server.
-spec start_link(Opts :: proplists:proplist()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Set the gateway PID for message routing.
%% Called by supervisor after both quic_server and gateway have started.
-spec set_gateway(pid(), pid()) -> ok.
set_gateway(QuicServerPid, GatewayPid) when is_pid(QuicServerPid), is_pid(GatewayPid) ->
gen_server:call(QuicServerPid, {set_gateway, GatewayPid}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @doc Initialize the QUIC server and start QUIC listener.
init(Opts) ->
Port = proplists:get_value(port, Opts, 9443),
Realm = proplists:get_value(realm, Opts, <<"macula.default">>),
GatewayPid = proplists:get_value(gateway, Opts), % Passed by gateway
CertFile = proplists:get_value(cert_file, Opts),
KeyFile = proplists:get_value(key_file, Opts),
NodeId = get_node_id(Realm, Port),
?LOG_INFO("Initializing QUIC server for realm ~s on port ~p", [Realm, Port]),
%% Start QUIC listener
ListenOpts = [
{cert, CertFile},
{key, KeyFile},
{alpn, ["macula"]},
{peer_unidi_stream_count, 3},
{peer_bidi_stream_count, 100}
],
case macula_quic:listen(Port, ListenOpts) of
{ok, Listener} ->
?LOG_INFO("QUIC listener started on port ~p", [Port]),
%% Start async accept to receive connections
case quicer:async_accept(Listener, #{}) of
{ok, Listener} ->
?LOG_INFO("Async accept registered", []),
ok;
{error, AcceptErr} ->
?LOG_WARNING("async_accept failed: ~p", [AcceptErr]),
ok
end,
State = #state{
listener = Listener,
port = Port,
realm = Realm,
node_id = NodeId,
gateway = GatewayPid
},
{ok, State};
{error, ErrorType, ErrorDetail} ->
?LOG_ERROR("QUIC listen failed: ~p ~p", [ErrorType, ErrorDetail]),
?LOG_ERROR("Certificate file: ~p", [CertFile]),
?LOG_ERROR("Key file: ~p", [KeyFile]),
{stop, {listen_failed, {ErrorType, ErrorDetail}}};
{error, Reason} ->
?LOG_ERROR("QUIC listen failed: ~p", [Reason]),
{stop, {listen_failed, Reason}}
end.
%% @doc Handle synchronous calls.
%% Set gateway PID for message routing
handle_call({set_gateway, GatewayPid}, _From, State) when is_pid(GatewayPid) ->
?LOG_INFO("Gateway PID set: ~p", [GatewayPid]),
{reply, ok, State#state{gateway = GatewayPid}};
%% Unknown calls
handle_call(_Request, _From, State) ->
{reply, {error, unknown_call}, State}.
%% @doc Handle asynchronous casts.
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle QUIC event: new_stream (stream created by peer).
handle_info({quic, new_stream, Stream, StreamProps}, State) ->
?LOG_DEBUG("========================================", []),
?LOG_DEBUG("NEW STREAM RECEIVED!", []),
?LOG_DEBUG("Stream: ~p", [Stream]),
?LOG_DEBUG("StreamProps: ~p", [StreamProps]),
?LOG_DEBUG("========================================", []),
%% Set stream to active mode to receive data automatically
case quicer:setopt(Stream, active, true) of
ok ->
?LOG_DEBUG("Stream set to active mode", []),
{noreply, State};
{error, Reason} ->
?LOG_ERROR("Failed to set stream active: ~p", [Reason]),
{noreply, State}
end;
%% @doc Handle QUIC data reception - decode and route to gateway.
%% Pattern matches on binary data, uses buffer for partial messages.
handle_info({quic, Data, Stream, _Flags}, State) when is_binary(Data) ->
?LOG_DEBUG("===== Received ~p bytes from stream ~p =====",
[byte_size(Data), Stream]),
?LOG_DEBUG("Raw data (first 100 bytes): ~p",
[binary:part(Data, 0, min(100, byte_size(Data)))]),
%% Decode message and route to gateway
DecodeResult = macula_protocol_decoder:decode(Data),
route_to_gateway(DecodeResult, Stream, State);
%% @doc Handle QUIC event: new_conn (connection established).
handle_info({quic, new_conn, Conn, ConnInfo}, State) ->
?LOG_DEBUG("========================================", []),
?LOG_DEBUG("NEW CONNECTION RECEIVED!", []),
?LOG_DEBUG("Connection: ~p", [Conn]),
?LOG_DEBUG("Connection Info: ~p", [ConnInfo]),
?LOG_DEBUG("========================================", []),
%% Use helper functions to complete handshake
complete_handshake(Conn),
register_next_connection(State#state.listener),
{noreply, State};
%% @doc Handle QUIC event: peer_needs_streams (peer wants to open more streams).
handle_info({quic, peer_needs_streams, _Conn, _StreamType}, State) ->
%% Normal QUIC control event - just acknowledge
{noreply, State};
%% @doc Handle QUIC event: shutdown (connection shutting down).
handle_info({quic, shutdown, Conn, Reason}, State) ->
?LOG_INFO("QUIC shutdown: Conn=~p, Reason=~p", [Conn, Reason]),
{noreply, State};
%% @doc Handle QUIC event: transport_shutdown (transport layer shutting down).
handle_info({quic, transport_shutdown, Conn, Reason}, State) ->
?LOG_INFO("QUIC transport_shutdown: Conn=~p, Reason=~p", [Conn, Reason]),
{noreply, State};
%% @doc Handle unknown messages.
handle_info(_Info, State) ->
{noreply, State}.
%% @doc Cleanup on termination.
terminate(_Reason, _State) ->
?LOG_INFO("Shutting down QUIC server", []),
ok.
%%%===================================================================
%%% Internal Helper Functions
%%%===================================================================
%% @doc Generate node ID from HOSTNAME env var (set by Docker) or generate from {Realm, Port}.
%% Returns a 32-byte binary (raw binary for Kademlia, never hex-encoded).
%% MUST match macula_gateway_system:get_node_id/2 and macula_gateway:get_node_id/2!
%%
%% Priority:
%% 1. NODE_NAME env var (explicit, highest priority)
%% 2. HOSTNAME env var (Docker sets this to container hostname - unique per container)
%% 3. Fallback to {Realm, Port} only (NO MAC - MAC is shared across Docker containers)
get_node_id(Realm, Port) when is_binary(Realm), is_integer(Port) ->
case os:getenv("NODE_NAME") of
false ->
%% No NODE_NAME, try HOSTNAME (Docker sets this to container hostname)
case os:getenv("HOSTNAME") of
false ->
%% No HOSTNAME either, use {Realm, Port} as last resort
crypto:hash(sha256, term_to_binary({Realm, Port}));
Hostname when is_list(Hostname) ->
%% Use HOSTNAME from Docker - unique per container
crypto:hash(sha256, term_to_binary({Realm, list_to_binary(Hostname), Port}))
end;
NodeName when is_list(NodeName) ->
%% Use NODE_NAME from environment - hash it to get 32-byte binary
crypto:hash(sha256, list_to_binary(NodeName))
end.
%%%===================================================================
%%% QUIC Connection Helper Functions
%%%===================================================================
%% @doc Complete TLS handshake on QUIC connection.
-spec complete_handshake(quicer:connection_handle()) -> ok.
complete_handshake(Conn) ->
case quicer:handshake(Conn) of
ok ->
accept_streams(Conn);
{ok, _} ->
accept_streams(Conn);
{error, Reason} ->
?LOG_ERROR("Handshake failed: ~p", [Reason]),
ok
end.
%% @doc Start accepting streams on an established connection.
-spec accept_streams(quicer:connection_handle()) -> ok.
accept_streams(Conn) ->
?LOG_INFO("Connection handshake completed successfully", []),
case quicer:async_accept_stream(Conn, #{}) of
{ok, Conn} ->
?LOG_INFO("Ready to accept streams on connection", []),
ok;
{error, StreamAcceptErr} ->
?LOG_WARNING("async_accept_stream failed: ~p", [StreamAcceptErr]),
ok
end.
%% @doc Register listener for next incoming connection.
-spec register_next_connection(quicer:listener_handle()) -> ok.
register_next_connection(Listener) ->
case quicer:async_accept(Listener, #{}) of
{ok, _} ->
?LOG_INFO("Ready for next connection", []),
ok;
{error, AcceptErr} ->
?LOG_WARNING("async_accept failed: ~p", [AcceptErr]),
ok
end.
%%%===================================================================
%%% Message Routing Functions
%%%===================================================================
%% @doc Route decoded message to gateway for business logic handling.
%% Pattern matches on decode result - uses declarative style.
-spec route_to_gateway(DecodeResult, Stream, State) -> {noreply, State}
when DecodeResult :: {ok, {atom(), map()}} | {error, term()},
Stream :: quicer:stream_handle(),
State :: #state{}.
%% Pattern 1: Successfully decoded message - route to gateway
route_to_gateway({ok, {MessageType, Message}}, Stream, State) when State#state.gateway =/= undefined ->
?LOG_DEBUG("Decoded message type: ~p", [MessageType]),
Gateway = State#state.gateway,
%% Route to gateway via gen_server:cast (async) to prevent blocking
%% This avoids timeout issues when gateway is busy processing other messages
gen_server:cast(Gateway, {route_message, MessageType, Message, Stream}),
?LOG_DEBUG("Message routed (async)", []),
{noreply, State};
%% Pattern 2: No gateway configured yet - log warning
route_to_gateway({ok, {MessageType, _Message}}, _Stream, State) ->
?LOG_WARNING("No gateway configured, dropping message type: ~p",
[MessageType]),
{noreply, State};
%% Pattern 3: Decode error - log and continue
route_to_gateway({error, Reason}, _Stream, State) ->
?LOG_ERROR("Decode error: ~p", [Reason]),
{noreply, State}.
%%%===================================================================
%%% Endpoint Parsing Functions
%%%===================================================================
%% @doc Parse endpoint URL to address tuple.
%% Converts "https://host:port" to {{IP_tuple}, Port}
%% Returns placeholder on parsing errors instead of crashing.
-spec parse_endpoint(undefined | binary()) -> {{byte(), byte(), byte(), byte()}, inet:port_number()}.
parse_endpoint(undefined) ->
{{0,0,0,0}, 0};
parse_endpoint(Endpoint) when is_binary(Endpoint) ->
%% Parse URL using uri_string
case uri_string:parse(Endpoint) of
#{host := Host, port := Port} when is_integer(Port) ->
resolve_host(Host, Port);
#{host := Host} ->
resolve_host(Host, 9443); %% Default port
_ ->
?LOG_WARNING("Invalid endpoint URL format: ~s, using placeholder", [Endpoint]),
{{0,0,0,0}, 0}
end.
%% @doc Resolve hostname to IP address.
%% Returns localhost fallback on DNS resolution failure.
-spec resolve_host(binary(), inet:port_number()) -> {{byte(), byte(), byte(), byte()}, inet:port_number()}.
resolve_host(Host, Port) when is_binary(Host), is_integer(Port) ->
HostStr = binary_to_list(Host),
case inet:getaddr(HostStr, inet) of
{ok, IPTuple} ->
{IPTuple, Port};
{error, Reason} ->
?LOG_WARNING("Failed to resolve host ~s: ~p, using localhost fallback",
[Host, Reason]),
{{127,0,0,1}, Port}
end.