Packages

macula

0.20.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_quic_stream_acceptor.erl
Raw

src/macula_quic_stream_acceptor.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% QUIC stream acceptor process.
%%% Dedicated process that waits for incoming streams on a connection
%%% and forwards them to the gateway for processing.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_quic_stream_acceptor).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([start_link/2]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-record(state, {
gateway_pid :: pid(),
conn :: term(),
accepting = false :: boolean()
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start stream acceptor process
-spec start_link(pid(), term()) -> {ok, pid()} | {error, term()}.
start_link(GatewayPid, Conn) ->
gen_server:start_link(?MODULE, [GatewayPid, Conn], []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init([GatewayPid, Conn]) ->
?LOG_INFO("Starting for connection ~p", [Conn]),
?LOG_DEBUG("Gateway PID: ~p", [GatewayPid]),
%% Register for stream events with active mode
StreamOpts = #{active => true},
case quicer:async_accept_stream(Conn, StreamOpts) of
{ok, Conn} ->
?LOG_DEBUG("Registered for streams (active mode)"),
{ok, #state{
gateway_pid = GatewayPid,
conn = Conn,
accepting = true
}};
{error, Reason} ->
?LOG_ERROR("Failed to register: ~p", [Reason]),
{stop, Reason}
end.
handle_call(Request, _From, State) ->
?LOG_WARNING("Unexpected call: ~p", [Request]),
{reply, {error, unknown_call}, State}.
handle_cast(Msg, State) ->
?LOG_WARNING("Unexpected cast: ~p", [Msg]),
{noreply, State}.
%% @doc Handle new stream from peer - THIS IS THE KEY MESSAGE!
handle_info({quic, new_stream, Stream, Props}, #state{gateway_pid = GatewayPid, conn = Conn} = State) ->
?LOG_DEBUG("========================================"),
?LOG_INFO("NEW STREAM RECEIVED!"),
?LOG_DEBUG("Stream: ~p", [Stream]),
?LOG_DEBUG("Props: ~p", [Props]),
?LOG_DEBUG("========================================"),
%% Forward stream to gateway
GatewayPid ! {quic_stream, Stream, Props},
%% Register for next stream
StreamOpts = #{active => true},
case quicer:async_accept_stream(Conn, StreamOpts) of
{ok, Conn} ->
?LOG_DEBUG("Re-registered for next stream"),
{noreply, State};
{error, Reason} ->
?LOG_ERROR("Failed to re-register: ~p", [Reason]),
{stop, {re_register_failed, Reason}, State}
end;
%% Handle stream data (if stream is in active mode)
handle_info({quic, Data, Stream, Flags}, #state{gateway_pid = GatewayPid} = State) when is_binary(Data) ->
?LOG_DEBUG("Stream data: Stream=~p, Size=~p bytes",
[Stream, byte_size(Data)]),
%% Forward to gateway
GatewayPid ! {quic_stream_data, Stream, Data, Flags},
{noreply, State};
%% Handle stream closed
handle_info({quic, stream_closed, Stream, Flags}, State) ->
?LOG_INFO("Stream closed: ~p, Flags: ~p", [Stream, Flags]),
{noreply, State};
%% Handle connection closed
handle_info({quic, closed, Conn, _Flags}, #state{conn = Conn} = State) ->
?LOG_INFO("Connection closed: ~p", [Conn]),
{stop, normal, State};
handle_info(Info, State) ->
?LOG_WARNING("Unhandled message: ~p", [Info]),
{noreply, State}.
terminate(Reason, #state{conn = Conn}) ->
?LOG_INFO("Terminating: ~p (Conn: ~p)", [Reason, Conn]),
ok.