Packages
macula
0.8.3
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_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).
%% 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),
io:format("[QuicServer] Initializing QUIC server for realm ~s on port ~p~n", [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} ->
io:format("[QuicServer] QUIC listener started on port ~p~n", [Port]),
%% Start async accept to receive connections
case quicer:async_accept(Listener, #{}) of
{ok, Listener} ->
io:format("[QuicServer] Async accept registered~n"),
ok;
{error, AcceptErr} ->
io:format("[QuicServer] WARNING: async_accept failed: ~p~n", [AcceptErr]),
ok
end,
State = #state{
listener = Listener,
port = Port,
realm = Realm,
node_id = NodeId,
gateway = GatewayPid
},
{ok, State};
{error, ErrorType, ErrorDetail} ->
io:format("[QuicServer] QUIC listen failed: ~p ~p~n", [ErrorType, ErrorDetail]),
io:format("[QuicServer] Certificate file: ~p~n", [CertFile]),
io:format("[QuicServer] Key file: ~p~n", [KeyFile]),
{stop, {listen_failed, {ErrorType, ErrorDetail}}};
{error, Reason} ->
io:format("[QuicServer] QUIC listen failed: ~p~n", [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) ->
io:format("[QuicServer] Gateway PID set: ~p~n", [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) ->
io:format("[QuicServer] ========================================~n"),
io:format("[QuicServer] NEW STREAM RECEIVED!~n"),
io:format("[QuicServer] Stream: ~p~n", [Stream]),
io:format("[QuicServer] StreamProps: ~p~n", [StreamProps]),
io:format("[QuicServer] ========================================~n"),
%% Set stream to active mode to receive data automatically
case quicer:setopt(Stream, active, true) of
ok ->
io:format("[QuicServer] Stream set to active mode~n"),
{noreply, State};
{error, Reason} ->
io:format("[QuicServer] Failed to set stream active: ~p~n", [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) ->
io:format("[QuicServer] ===== Received ~p bytes from stream ~p =====~n",
[byte_size(Data), Stream]),
io:format("[QuicServer] Raw data (first 100 bytes): ~p~n",
[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) ->
io:format("[QuicServer] ========================================~n"),
io:format("[QuicServer] NEW CONNECTION RECEIVED!~n"),
io:format("[QuicServer] Connection: ~p~n", [Conn]),
io:format("[QuicServer] Connection Info: ~p~n", [ConnInfo]),
io:format("[QuicServer] ========================================~n"),
%% 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) ->
io:format("[QuicServer] QUIC shutdown: Conn=~p, Reason=~p~n", [Conn, Reason]),
{noreply, State};
%% @doc Handle QUIC event: transport_shutdown (transport layer shutting down).
handle_info({quic, transport_shutdown, Conn, Reason}, State) ->
io:format("[QuicServer] QUIC transport_shutdown: Conn=~p, Reason=~p~n", [Conn, Reason]),
{noreply, State};
%% @doc Handle unknown messages.
handle_info(_Info, State) ->
{noreply, State}.
%% @doc Cleanup on termination.
terminate(_Reason, _State) ->
io:format("[QuicServer] Shutting down QUIC server~n"),
ok.
%%%===================================================================
%%% Internal Helper Functions
%%%===================================================================
%% @doc Generate node ID from realm and port.
%% Uses MD5 hash of "realm:port" similar to gateway implementation.
get_node_id(Realm, Port) when is_binary(Realm), is_integer(Port) ->
PortBin = integer_to_binary(Port),
Input = <<Realm/binary, ":", PortBin/binary>>,
crypto:hash(md5, Input).
%%%===================================================================
%%% 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} ->
io:format("[QuicServer] Handshake failed: ~p~n", [Reason]),
ok
end.
%% @doc Start accepting streams on an established connection.
-spec accept_streams(quicer:connection_handle()) -> ok.
accept_streams(Conn) ->
io:format("[QuicServer] Connection handshake completed successfully~n"),
case quicer:async_accept_stream(Conn, #{}) of
{ok, Conn} ->
io:format("[QuicServer] Ready to accept streams on connection~n"),
ok;
{error, StreamAcceptErr} ->
io:format("[QuicServer] WARNING: async_accept_stream failed: ~p~n", [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, _} ->
io:format("[QuicServer] Ready for next connection~n"),
ok;
{error, AcceptErr} ->
io:format("[QuicServer] WARNING: async_accept failed: ~p~n", [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 ->
io:format("[QuicServer] Decoded message type: ~p~n", [MessageType]),
Gateway = State#state.gateway,
%% Route to gateway via gen_server:call with extended timeout
%% (mesh connection creation can take longer than default 5s)
case gen_server:call(Gateway, {route_message, MessageType, Message, Stream}, 30000) of
ok ->
io:format("[QuicServer] Message routed successfully~n"),
{noreply, State};
{error, Reason} ->
io:format("[QuicServer] Routing failed: ~p~n", [Reason]),
{noreply, State}
end;
%% Pattern 2: No gateway configured yet - log warning
route_to_gateway({ok, {MessageType, _Message}}, _Stream, State) ->
io:format("[QuicServer] WARNING: No gateway configured, dropping message type: ~p~n",
[MessageType]),
{noreply, State};
%% Pattern 3: Decode error - log and continue
route_to_gateway({error, Reason}, _Stream, State) ->
io:format("[QuicServer] Decode error: ~p~n", [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
_ ->
io:format("[QuicServer] Invalid endpoint URL format: ~s, using placeholder~n", [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} ->
io:format("[QuicServer] Failed to resolve host ~s: ~p, using localhost fallback~n",
[Host, Reason]),
{{127,0,0,1}, Port}
end.