Packages

macula

4.2.4
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_net macula_net_transport_quic.erl
Raw

src/macula_net/macula_net_transport_quic.erl

%%%-------------------------------------------------------------------
%%% @doc QUIC transport plugin for macula-net.
%%%
%%% Implements {@link macula_net_transport} using the SDK's existing
%%% {@link macula_quic} primitives (Quinn-based Rust NIF). One bidi
%%% stream per outbound connection carries length-prefixed CBOR
%%% envelopes:
%%%
%%% <pre>
%%% &lt;&lt;Len:32/big, Cbor:Len/binary&gt;&gt;
%%% </pre>
%%%
%%% == Phase 1 simplifications ==
%%%
%%% These will be addressed in Phase 4 hardening (PLAN_MACULA_NET.md §13):
%%% <ul>
%%% <li>Self-signed throwaway TLS cert generated at startup. Identity
%%% authentication happens at the macula-net envelope layer (sigs
%%% on control messages), NOT at the TLS layer. Phase 4 swaps to
%%% raw-pubkey TLS bound to the macula identity.</li>
%%% <li>One outbound connection per peer; no connection pooling.</li>
%%% <li>On disconnect, no automatic reconnect.</li>
%%% <li>No backpressure beyond the QUIC flow-control window.</li>
%%% </ul>
%%% @end
%%%-------------------------------------------------------------------
-module(macula_net_transport_quic).
-behaviour(macula_net_transport).
-behaviour(gen_server).
-export([
start_link/1,
stop/0,
set_handler/1,
connect/3,
disconnect/1,
send/2,
peer_path_mtu/1,
connection_count/0
]).
%% gen_server
-export([init/1, handle_call/3, handle_cast/2, handle_info/2,
terminate/2, code_change/3]).
-define(SERVER, ?MODULE).
-define(FRAME_HEADER_BYTES, 4).
-define(MAX_FRAME_BYTES, 16#100000). %% 1 MiB cap
-define(STREAM_OPEN_TIMEOUT_MS, 5000).
-define(CONNECT_TIMEOUT_MS, 5000).
-define(PATH_MTU_TICK_MS, 5000).
-record(out_link, {
conn :: reference(),
stream :: reference()
}).
-record(in_state, {
%% Inbound stream framing buffer (per stream).
buf = <<>> :: binary()
}).
-record(state, {
listener :: reference() | undefined,
listen_port :: inet:port_number() | undefined,
cert_path :: string() | undefined,
key_path :: string() | undefined,
handler :: macula_net_transport:handler() | undefined,
%% station_id() => #out_link{}
out = #{} :: #{macula_net_transport:station_id() => #out_link{}},
%% inbound stream ref => #in_state{}
in_streams = #{}:: #{reference() => #in_state{}}
}).
%% =============================================================================
%% Public API
%% =============================================================================
%% @doc Start the QUIC listener.
%%
%% `Opts' map keys:
%% <ul>
%% <li>`port' — UDP port for the QUIC listener (mandatory)</li>
%% <li>`bind' — bind address binary (default `&lt;&lt;"::"&gt;&gt;')</li>
%% <li>`alpn' — ALPN protocol id binary (default `&lt;&lt;"macula-net"&gt;&gt;')</li>
%% </ul>
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(#{port := _} = Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [Opts], []).
-spec stop() -> ok.
stop() ->
case whereis(?SERVER) of
undefined -> ok;
_ -> gen_server:stop(?SERVER)
end.
%% @doc Register the inbound handler. Called for every received envelope.
%%
%% The handler is invoked as `Handler(Cbor, StreamRef)' where StreamRef
%% is the bidi-stream reference the frame arrived on. Non-host handlers
%% can ignore StreamRef; the host_attach_controller (Phase 3.5) uses it
%% to forward replies on the same stream a daemon dialed in on.
-spec set_handler(macula_net_transport:handler()) -> ok.
set_handler(Handler) when is_function(Handler, 2) ->
gen_server:call(?SERVER, {set_handler, Handler}).
%% @doc Open an outbound QUIC connection + bidi stream to a peer station.
-spec connect(StationId :: macula_net_transport:station_id(),
Host :: binary() | string(),
Port :: inet:port_number()) -> ok | {error, term()}.
connect(StationId, Host, Port) ->
gen_server:call(?SERVER, {connect, StationId, Host, Port},
?CONNECT_TIMEOUT_MS + 1000).
-spec disconnect(macula_net_transport:station_id()) -> ok.
disconnect(StationId) ->
gen_server:call(?SERVER, {disconnect, StationId}).
%% @doc Path MTU as currently tracked by Quinn for the connection to
%% `StationId'. Phase 4.2 — see PLAN_MACULA_NET_PHASE4_2_MTU_PMTUD.md.
-spec peer_path_mtu(macula_net_transport:station_id()) ->
{ok, pos_integer()} | {error, not_connected | term()}.
peer_path_mtu(StationId) ->
gen_server:call(?SERVER, {peer_path_mtu, StationId}).
%% @doc Active outbound-connection count. Used by macula_metrics's
%% gauge poller (see PLAN_MACULA_NET_PHASE4_1_OBSERVABILITY.md).
-spec connection_count() -> non_neg_integer().
connection_count() ->
case whereis(?SERVER) of
undefined -> 0;
_ -> gen_server:call(?SERVER, connection_count)
end.
%% @doc Send a CBOR envelope to a known station (must be `connect'ed first).
-spec send(macula_net_transport:station_id(),
macula_net_transport:cbor_envelope()) -> ok | {error, term()}.
send(StationId, Cbor) ->
gen_server:call(?SERVER, {send, StationId, Cbor}).
%% =============================================================================
%% gen_server callbacks
%% =============================================================================
init([#{port := Port} = Opts]) ->
process_flag(trap_exit, true),
BindAddr = maps:get(bind, Opts, <<"::">>),
Alpn = maps:get(alpn, Opts, <<"macula-net">>),
{ok, {CertPath, KeyPath}} = ensure_self_signed_cert(Port),
ListenOpts = [
{cert, CertPath},
{key, KeyPath},
{alpn, [Alpn]},
{idle_timeout_ms, 120000},
{keep_alive_interval_ms, 30000}
],
case macula_quic:listen(BindAddr, Port, ListenOpts) of
{ok, Listener} ->
ok = macula_quic:async_accept(Listener),
erlang:send_after(?PATH_MTU_TICK_MS, self(), tick_path_mtu),
{ok, #state{listener = Listener,
listen_port = Port,
cert_path = CertPath,
key_path = KeyPath}};
{error, Reason} ->
{stop, Reason}
end.
handle_call({set_handler, Handler}, _From, State) ->
{reply, ok, State#state{handler = Handler}};
handle_call({connect, StationId, Host, Port}, _From, #state{out = Out} = State) ->
handle_connect(maps:is_key(StationId, Out), StationId, Host, Port, State);
handle_call({disconnect, StationId}, _From, #state{out = Out} = State) ->
{Reply, NewState} = drop_out_link(StationId, Out, State),
{reply, Reply, NewState};
handle_call({send, StationId, Cbor}, _From, #state{out = Out} = State) ->
{reply, do_send(maps:get(StationId, Out, undefined), Cbor), State};
handle_call({peer_path_mtu, StationId}, _From, #state{out = Out} = State) ->
{reply, do_peer_path_mtu(maps:get(StationId, Out, undefined)), State};
handle_call(connection_count, _From, #state{out = Out} = State) ->
{reply, map_size(Out), State};
handle_call(_Other, _From, State) ->
{reply, {error, unknown_call}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% A new inbound QUIC connection arrived on the listener.
handle_info({quic, new_conn, Conn, _ConnInfo},
#state{listener = Listener} = State) ->
%% Hand off control to ourselves; arm the next accept.
ok = macula_quic:controlling_process(Conn, self()),
ok = macula_quic:async_accept_stream(Conn),
ok = macula_quic:async_accept(Listener),
{noreply, State};
%% A new bidi stream arrived on an inbound connection.
handle_info({quic, new_stream, Stream, _StreamInfo},
#state{in_streams = InStreams} = State) ->
ok = macula_quic:controlling_process(Stream, self()),
ok = macula_quic:setopt(Stream, active, true),
telemetry:execute([macula, net, transport, stream_opened],
#{count => 1}, #{direction => <<"inbound">>}),
{noreply, State#state{in_streams = InStreams#{Stream => #in_state{}}}};
%% Bytes arrived on an inbound stream. macula_quic delivers data as
%% `{quic, Binary, StreamRef, Flags}' — note the second element is
%% the binary itself, NOT a `data' tag. See
%% native/macula_quic/src/message.rs.
handle_info({quic, Data, Stream, _Flags}, State)
when is_binary(Data), is_reference(Stream) ->
{noreply, deliver_buffered(Stream, Data, State)};
%% Stream / connection lifecycle: clean up tracking maps.
handle_info({quic, stream_closed, Stream, _Reason},
#state{in_streams = InStreams} = State) ->
telemetry:execute([macula, net, transport, stream_closed],
#{count => 1}, #{direction => <<"inbound">>}),
{noreply, State#state{in_streams = maps:remove(Stream, InStreams)}};
handle_info({quic, conn_closed, _Conn, _Reason}, State) ->
{noreply, State};
handle_info({'EXIT', _Pid, _Reason}, State) ->
{noreply, State};
handle_info(tick_path_mtu, #state{out = Out} = State) ->
maps:foreach(fun emit_path_mtu/2, Out),
erlang:send_after(?PATH_MTU_TICK_MS, self(), tick_path_mtu),
{noreply, State};
handle_info(_Other, State) ->
{noreply, State}.
terminate(_Reason, #state{listener = L, out = Out, cert_path = CP, key_path = KP}) ->
case L of
undefined -> ok;
_ -> _ = macula_quic:close_listener(L), ok
end,
maps:foreach(
fun(_K, #out_link{conn = Conn}) -> _ = macula_quic:close_connection(Conn), ok end,
Out),
%% Clean up the throwaway cert/key files.
case CP of
undefined -> ok;
_ -> _ = file:delete(CP), ok
end,
case KP of
undefined -> ok;
_ -> _ = file:delete(KP), ok
end,
ok.
code_change(_OldVsn, State, _Extra) -> {ok, State}.
%% =============================================================================
%% Connect / disconnect helpers
%% =============================================================================
handle_connect(true, _StationId, _Host, _Port, State) ->
telemetry:execute([macula, net, transport, connect],
#{count => 1}, #{outcome => <<"already_connected">>}),
{reply, {error, already_connected}, State};
handle_connect(false, StationId, Host, Port, State) ->
after_quic_connect(macula_quic:connect(to_binary(Host), Port,
default_connect_opts(),
?CONNECT_TIMEOUT_MS),
StationId, State).
after_quic_connect({error, Reason} = E, _StationId, State) ->
telemetry:execute([macula, net, transport, connect],
#{count => 1},
#{outcome => connect_outcome(Reason)}),
{reply, E, State};
after_quic_connect({ok, Conn}, StationId, State) ->
ok = macula_quic:controlling_process(Conn, self()),
after_open_stream(macula_quic:open_stream(Conn), Conn, StationId, State).
after_open_stream({error, Reason} = E, Conn, _StationId, State) ->
_ = macula_quic:close_connection(Conn),
telemetry:execute([macula, net, transport, connect],
#{count => 1},
#{outcome => connect_outcome(Reason)}),
{reply, E, State};
after_open_stream({ok, Stream}, Conn, StationId,
#state{out = Out, in_streams = InStreams} = State) ->
%% Phase 3.5 — activate recv on outbound streams so the host can
%% send forwarded data back on the same bidi stream the daemon
%% dialed in on. Phase 1/2 didn't surface this because traffic
%% only flowed one way (Bob -> Helsinki -> TUN); 3.6 needs the
%% return path. Track the stream in in_streams so deliver_buffered
%% has a framing buffer for it.
ok = macula_quic:setopt(Stream, active, true),
Link = #out_link{conn = Conn, stream = Stream},
telemetry:execute([macula, net, transport, connect],
#{count => 1}, #{outcome => <<"ok">>}),
telemetry:execute([macula, net, transport, stream_opened],
#{count => 1}, #{direction => <<"outbound">>}),
{reply, ok, State#state{
out = Out#{StationId => Link},
in_streams = InStreams#{Stream => #in_state{}}
}}.
connect_outcome(timeout) -> <<"timeout">>;
connect_outcome(connection_refused) -> <<"refused">>;
connect_outcome(R) when is_atom(R) -> atom_to_binary(R, utf8);
connect_outcome(R) when is_binary(R) -> R;
connect_outcome(_) -> <<"error">>.
drop_out_link(StationId, Out, State) ->
case maps:take(StationId, Out) of
{#out_link{conn = Conn, stream = Stream}, NewOut} ->
_ = macula_quic:close_stream(Stream),
_ = macula_quic:close_connection(Conn),
{ok, State#state{out = NewOut}};
error ->
{ok, State}
end.
do_send(undefined, _Cbor) ->
{error, not_connected};
do_send(#out_link{stream = Stream}, Cbor) ->
Frame = <<(byte_size(Cbor)):32/big, Cbor/binary>>,
macula_quic:send(Stream, Frame).
do_peer_path_mtu(undefined) ->
{error, not_connected};
do_peer_path_mtu(#out_link{conn = Conn}) ->
macula_quic:max_datagram_size(Conn).
emit_path_mtu(StationId, #out_link{conn = Conn}) ->
decide_emit_mtu(macula_quic:max_datagram_size(Conn), StationId).
decide_emit_mtu({ok, Bytes}, StationId) ->
telemetry:execute([macula, net, transport, path_mtu],
#{bytes => Bytes},
#{peer => station_label(StationId)});
decide_emit_mtu({error, _}, _StationId) ->
ok.
station_label(StationId) when is_binary(StationId), byte_size(StationId) >= 8 ->
<<Prefix:8/binary, _/binary>> = StationId,
bin_to_hex(Prefix);
station_label(StationId) when is_binary(StationId) ->
bin_to_hex(StationId);
station_label(_) ->
<<"unknown">>.
bin_to_hex(Bin) ->
list_to_binary([io_lib:format("~2.16.0b", [B]) || <<B>> <= Bin]).
%% =============================================================================
%% Inbound framing
%% =============================================================================
deliver_buffered(Stream, NewBytes, #state{handler = Handler,
in_streams = InStreams} = State) ->
Existing = maps:get(Stream, InStreams, #in_state{}),
Combined = <<(Existing#in_state.buf)/binary, NewBytes/binary>>,
{Frames, Rest} = extract_frames(Combined, []),
lists:foreach(fun(F) -> safe_invoke(Handler, F, Stream) end, Frames),
State#state{in_streams = InStreams#{Stream => #in_state{buf = Rest}}}.
extract_frames(<<Len:32/big, Cbor:Len/binary, Rest/binary>>, Acc)
when Len =< ?MAX_FRAME_BYTES ->
extract_frames(Rest, [Cbor | Acc]);
extract_frames(<<Len:32/big, _/binary>> = _Bin, Acc) when Len > ?MAX_FRAME_BYTES ->
%% Oversized frame — drop the whole buffer; recovery is connection
%% reset territory. Phase 4 hardening.
{lists:reverse(Acc), <<>>};
extract_frames(Bin, Acc) ->
{lists:reverse(Acc), Bin}.
safe_invoke(undefined, _Cbor, _Stream) -> ok;
safe_invoke(Handler, Cbor, Stream) when is_function(Handler, 2) ->
%% Per CLAUDE.md: avoid try/catch where possible. We accept it here
%% because a user-supplied handler crashing must not take down the
%% transport. This is the boundary; below it we let it crash.
try Handler(Cbor, Stream)
catch _:_ -> ok
end.
%% =============================================================================
%% Utilities
%% =============================================================================
to_binary(X) when is_binary(X) -> X;
to_binary(X) when is_list(X) -> list_to_binary(X).
default_connect_opts() ->
[
{alpn, [<<"macula-net">>]},
%% Phase 1: skip server name verification — self-signed certs.
%% Phase 4 hardening swaps to raw-pubkey TLS pinned to identity.
{verify, none},
{idle_timeout_ms, 120000},
{keep_alive_interval_ms, 30000}
].
%% Generate a throwaway Ed25519 keypair + self-signed cert at boot.
%% Phase 4 swaps this for identity-bound certs.
ensure_self_signed_cert(Port) ->
{Pub, Priv} = ephemeral_ed25519_keypair(),
Sans = [<<"localhost">>, <<"::1">>, <<"127.0.0.1">>],
{ok, {CertPem, KeyPem}} = macula_quic:generate_self_signed_cert(Pub, Priv, Sans),
Tmp = case os:getenv("TMPDIR") of
false -> "/tmp";
V -> V
end,
%% Include the OS pid in the stem so multiple BEAM nodes sharing
%% the same /tmp (e.g. several netns on one box) don't race on
%% writing the cert/key files. Without this, Phase 2's three-netns
%% demo hits "TLS keys may not be consistent: KeyMismatch".
Stem = filename:join(Tmp, io_lib:format("macula-net-~p-~s",
[Port, os:getpid()])),
CertPath = lists:flatten([Stem, ".crt"]),
KeyPath = lists:flatten([Stem, ".key"]),
ok = file:write_file(CertPath, CertPem),
ok = file:write_file(KeyPath, KeyPem),
{ok, {CertPath, KeyPath}}.
ephemeral_ed25519_keypair() ->
%% crypto:generate_key/2 uses the BEAM's OpenSSL; suitable for a
%% throwaway TLS cert in Phase 1.
{Pub, Priv} = crypto:generate_key(eddsa, ed25519),
{iolist_to_binary(Pub), iolist_to_binary(Priv)}.