Packages

macula

3.14.0
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 peering macula_quic.erl
Raw

src/peering/macula_quic.erl

%%%-------------------------------------------------------------------
%%% @doc Macula QUIC transport — Quinn-based Rust NIF.
%%%
%%% Provides QUIC listener, connection, and stream operations backed
%%% by Quinn (Rust) instead of MsQuic. Key improvement: listeners
%%% actually bind to specific IP addresses, enabling per-identity
%%% IPv6 binding for virtual relay identities.
%%%
%%% Active-mode messages delivered to owning process:
%%% {quic, Data, StreamRef, Flags} — stream data
%%% {quic, new_conn, ConnRef, Info} — new connection accepted
%%% {quic, new_stream, StreamRef, Props} — new stream accepted
%%% {quic, peer_send_shutdown, StreamRef, undefined}
%%% {quic, stream_closed, StreamRef, Flags}
%%% {quic, shutdown, Handle, Reason}
%%% @end
%%%-------------------------------------------------------------------
-module(macula_quic).
-include_lib("kernel/include/logger.hrl").
-on_load(init/0).
-export([
%% Listener
listen/2,
listen/3,
async_accept/1,
async_accept/2,
close_listener/1,
%% Connection
connect/4,
open_stream/1,
%% Sovereign overlay (Yggdrasil) — self-signed cert generation
generate_self_signed_cert/3,
close_connection/1,
async_accept_stream/1,
async_accept_stream/2,
handshake/1,
peername/1,
%% Stream
send/2,
async_send/2,
close_stream/1,
setopt/3,
controlling_process/2,
%% Compat: generic close (tries stream → conn → listener)
close/1,
%% Dist compat (stream accept with opts, stream open with opts)
accept_stream/3,
open_stream/2,
handoff_stream/3,
%% Shutdown (maps to close with flags)
async_shutdown_stream/3,
async_shutdown_connection/3,
%% Stats
getstat/2
]).
%%%===================================================================
%%% NIF Loading
%%%===================================================================
init() ->
PrivDir = code:priv_dir(macula),
SoName = filename:join(PrivDir, "libmacula_quic"),
case erlang:load_nif(SoName, 0) of
ok ->
?LOG_INFO("[macula_quic] Quinn NIF loaded from ~s", [SoName]),
ok;
{error, {reload, _}} ->
ok;
{error, Reason} ->
?LOG_WARNING("[macula_quic] NIF load failed: ~p (path: ~s)", [Reason, SoName]),
{error, Reason}
end.
%%%===================================================================
%%% Listener API
%%%===================================================================
%% @doc Listen on a port or {Address, Port} tuple.
-spec listen(inet:port_number() | {string() | binary(), inet:port_number()}, list()) ->
{ok, reference()} | {error, term()}.
listen({Address, Port}, Opts) ->
listen(Address, Port, Opts);
listen(Port, Opts) when is_integer(Port) ->
listen(<<"::">>, Port, Opts).
%% @doc Listen on a specific bind address and port.
%% BindAddr is a binary: "0.0.0.0", "192.168.1.1", "2600:3c0e::100", etc.
-spec listen(binary() | string(), inet:port_number(), list()) -> {ok, reference()} | {error, term()}.
listen(BindAddr, Port, Opts) when is_list(BindAddr) ->
listen(list_to_binary(BindAddr), Port, Opts);
listen(BindAddr, Port, Opts) when is_binary(BindAddr) ->
CertFile = to_binary(proplists:get_value(cert, Opts)),
KeyFile = to_binary(proplists:get_value(key, Opts)),
Alpn = [to_binary(A) || A <- proplists:get_value(alpn, Opts, ["macula"])],
IdleTimeoutMs = proplists:get_value(idle_timeout_ms, Opts, 120000),
KeepAliveMs = proplists:get_value(keep_alive_interval_ms, Opts, 30000),
BidiStreams = proplists:get_value(peer_bidi_stream_count, Opts, 100),
UniStreams = proplists:get_value(peer_unidi_stream_count, Opts, 3),
?LOG_INFO("Starting listener on ~s:~p with idle_timeout=~pms, keep_alive=~pms",
[BindAddr, Port, IdleTimeoutMs, KeepAliveMs]),
nif_listen(BindAddr, Port, CertFile, KeyFile, Alpn,
IdleTimeoutMs, KeepAliveMs, BidiStreams, UniStreams).
%% @doc Start accepting connections on a listener.
%% Delivers {quic, new_conn, ConnRef, Info} to the calling process.
-spec async_accept(reference()) -> ok | {error, term()}.
async_accept(Listener) ->
async_accept(Listener, #{}).
-spec async_accept(reference(), map()) -> ok | {error, term()}.
async_accept(Listener, _Opts) ->
nif_async_accept(Listener).
%% @doc Close a listener.
-spec close_listener(reference()) -> ok.
close_listener(Listener) ->
nif_close_listener(Listener).
%%%===================================================================
%%% Connection API
%%%===================================================================
%% @doc Connect to a remote QUIC server.
%%
%% Target forms:
%% <ul>
%% <li>`Host :: binary() | string()' — the existing hostname /
%% IP-string path. Validation depends on `verify' / `verify_pubkey'
%% opts.</li>
%% <li>`{pubkey, Pubkey :: binary()}' — sovereign-overlay path.
%% The 32-byte Ed25519 pubkey is the identity. The Yggdrasil
%% IPv6 is derived from it; the leaf cert is validated by SPKI
%% pin against the same pubkey. No DNS, no CA. See
%% PLAN_SOVEREIGN_OVERLAY_PHASE1 §4.4.</li>
%% </ul>
-spec connect(Target, inet:port_number(), list(), timeout()) ->
{ok, reference()} | {error, term()}
when Target :: binary() | string() | {pubkey, binary()}.
connect({pubkey, Pubkey}, Port, Opts, Timeout)
when is_binary(Pubkey), byte_size(Pubkey) =:= 32 ->
Addr = macula_yggdrasil:address_for(Pubkey),
HostBin = <<"[", (macula_yggdrasil:format_address(Addr))/binary, "]">>,
%% Force pubkey-pin verification; override any conflicting opt.
Opts1 = lists:keystore(verify_pubkey, 1, Opts, {verify_pubkey, Pubkey}),
%% `verify' must stay falsy on this path — webpki would reject
%% the IP-form SNI server-name and refuse to load.
Opts2 = lists:keystore(verify, 1, Opts1, {verify, none}),
connect(HostBin, Port, Opts2, Timeout);
connect(Host, Port, Opts, Timeout) ->
HostBin = to_binary(Host),
Alpn = [to_binary(A) || A <- proplists:get_value(alpn, Opts, ["macula"])],
Verify = proplists:get_value(verify, Opts, none) =/= none,
%% `verify_pubkey' is a 32-byte Ed25519 pubkey to pin against the
%% leaf cert SPKI. Empty binary disables pinning. Sovereign
%% overlay path uses this to validate by pubkey alone (no CA).
VerifyPubkey = proplists:get_value(verify_pubkey, Opts, <<>>),
IdleTimeoutMs = proplists:get_value(idle_timeout_ms, Opts, 60000),
KeepAliveMs = proplists:get_value(keep_alive_interval_ms, Opts, 20000),
nif_connect(HostBin, Port, Alpn, Verify, VerifyPubkey,
IdleTimeoutMs, KeepAliveMs, Timeout).
%%%===================================================================
%%% Sovereign overlay (Yggdrasil) — self-signed cert generation
%%%===================================================================
%% @doc Generate a self-signed X.509 cert from an Ed25519 keypair.
%% Returns `{ok, {CertPem, KeyPem}}' as PEM-encoded binaries
%% suitable for handing to `macula_quic:listen/3' via
%% `cert' / `key' opts (after writing to disk).
%%
%% Used by `macula_yggdrasil:cert_for/1' for the sovereign-overlay
%% listener path — the cert wraps the identity's macula pubkey,
%% no CA chain required. See PLAN_SOVEREIGN_OVERLAY_PHASE1 §4.3.
-spec generate_self_signed_cert(Pubkey :: binary(),
Privkey :: binary(),
Sans :: [binary() | string()]) ->
{ok, {CertPem :: binary(), KeyPem :: binary()}} | {error, term()}.
generate_self_signed_cert(Pubkey, Privkey, Sans)
when is_binary(Pubkey), byte_size(Pubkey) =:= 32,
is_binary(Privkey), byte_size(Privkey) =:= 32 ->
SansCsv = iolist_to_binary(
lists:join(<<",">>, [to_binary(S) || S <- Sans])),
nif_generate_self_signed_cert(Pubkey, Privkey, SansCsv).
%% @doc Open a new bidirectional stream.
-spec open_stream(reference()) -> {ok, reference()} | {error, term()}.
open_stream(Conn) ->
nif_open_stream(Conn).
%% @doc Close a connection.
-spec close_connection(reference()) -> ok.
close_connection(Conn) ->
nif_close_connection(Conn).
%% @doc Start accepting streams on a connection.
%% Delivers {quic, new_stream, StreamRef, #{conn => ConnRef}} to the owning process.
-spec async_accept_stream(reference()) -> ok | {error, term()}.
async_accept_stream(Conn) ->
async_accept_stream(Conn, #{}).
-spec async_accept_stream(reference(), map()) -> ok | {error, term()}.
async_accept_stream(Conn, _Opts) ->
nif_async_accept_stream(Conn).
%% @doc Complete TLS handshake.
%% With Quinn, handshake completes during accept — this is a no-op for compat.
-spec handshake(reference()) -> ok | {ok, reference()} | {error, term()}.
handshake(Conn) ->
{ok, Conn}.
%% @doc Get remote address of a connection.
-spec peername(reference()) -> {ok, {string(), inet:port_number()}} | {error, term()}.
peername(Conn) ->
nif_peername(Conn).
%%%===================================================================
%%% Stream API
%%%===================================================================
%% @doc Send data on a stream (blocking).
-spec send(reference(), iodata()) -> ok | {error, term()}.
send(Stream, Data) ->
nif_send(Stream, iolist_to_binary(Data)).
%% @doc Send data asynchronously.
-spec async_send(reference(), iodata()) -> ok | {error, term()}.
async_send(Stream, Data) ->
nif_async_send(Stream, iolist_to_binary(Data)).
%% @doc Close a stream.
-spec close_stream(reference()) -> ok.
close_stream(Stream) ->
nif_close_stream(Stream).
%% @doc Set active mode on a stream handle.
-spec setopt(reference(), active, boolean()) -> ok.
setopt(Stream, active, Value) ->
nif_setopt_active(Stream, Value).
%% @doc Transfer ownership of a handle to another process.
%% Works with both stream and connection handles.
-spec controlling_process(reference(), pid()) -> ok | {error, term()}.
controlling_process(Handle, Pid) ->
try nif_controlling_process(Handle, Pid)
catch error:badarg ->
nif_controlling_process_conn(Handle, Pid)
end.
%%%===================================================================
%%% Compat API
%%%===================================================================
%% @doc Generic close — tries stream, then connection, then listener.
-spec close(reference()) -> ok.
close(Ref) ->
close_as(Ref, [fun nif_close_stream/1,
fun nif_close_connection/1,
fun nif_close_listener/1]).
close_as(_Ref, []) -> ok;
close_as(Ref, [CloseFn | Rest]) ->
try CloseFn(Ref) of ok -> ok
catch _:_ -> close_as(Ref, Rest)
end.
%% @doc Async shutdown stream (compat with quicer flags).
-spec async_shutdown_stream(reference(), integer(), integer()) -> ok.
async_shutdown_stream(Stream, _Flag, _Code) ->
nif_close_stream(Stream).
%% @doc Async shutdown connection (compat with quicer flags).
-spec async_shutdown_connection(reference(), integer(), integer()) -> ok.
async_shutdown_connection(Conn, _Flag, _Code) ->
nif_close_connection(Conn).
%% @doc Get connection stats. Currently returns zeros for all
%% requested counters — Quinn exposes per-connection stats but the
%% NIF binding hasn't surfaced them yet. Used by macula_dist for
%% telemetry; zeroed values are harmless (dist_util only uses these
%% for liveness signals).
-spec getstat(reference(), [atom()]) -> {ok, [{atom(), integer()}]} | {error, term()}.
getstat(_Conn, Stats) ->
{ok, [{S, 0} || S <- Stats]}.
%%%===================================================================
%%% Dist Compat API
%%%===================================================================
%% @doc Accept stream with options and timeout (for macula_dist).
-spec accept_stream(reference(), map(), timeout()) -> {ok, reference()} | {error, term()}.
accept_stream(Conn, _Opts, _Timeout) ->
async_accept_stream(Conn).
%% @doc Open stream with options map (for macula_dist).
-spec open_stream(reference(), map()) -> {ok, reference()} | {error, term()}.
open_stream(Conn, _Opts) ->
open_stream(Conn).
%% @doc Hand off a stream to another process (for macula_dist).
-spec handoff_stream(reference(), pid(), map()) -> ok | {error, term()}.
handoff_stream(Stream, NewOwner, _Opts) ->
controlling_process(Stream, NewOwner).
%%%===================================================================
%%% NIF Stubs
%%%===================================================================
nif_listen(_BindAddr, _Port, _CertFile, _KeyFile, _Alpn,
_IdleTimeoutMs, _KeepAliveMs, _BidiStreams, _UniStreams) ->
erlang:nif_error(nif_not_loaded).
nif_async_accept(_Listener) ->
erlang:nif_error(nif_not_loaded).
nif_close_listener(_Listener) ->
erlang:nif_error(nif_not_loaded).
nif_connect(_Host, _Port, _Alpn, _Verify, _VerifyPubkey,
_IdleTimeoutMs, _KeepAliveMs, _TimeoutMs) ->
erlang:nif_error(nif_not_loaded).
nif_generate_self_signed_cert(_Pubkey, _Privkey, _Sans) ->
erlang:nif_error(nif_not_loaded).
nif_open_stream(_Conn) ->
erlang:nif_error(nif_not_loaded).
nif_close_connection(_Conn) ->
erlang:nif_error(nif_not_loaded).
nif_async_accept_stream(_Conn) ->
erlang:nif_error(nif_not_loaded).
nif_peername(_Conn) ->
erlang:nif_error(nif_not_loaded).
nif_send(_Stream, _Data) ->
erlang:nif_error(nif_not_loaded).
nif_async_send(_Stream, _Data) ->
erlang:nif_error(nif_not_loaded).
nif_close_stream(_Stream) ->
erlang:nif_error(nif_not_loaded).
nif_setopt_active(_Stream, _Value) ->
erlang:nif_error(nif_not_loaded).
nif_controlling_process(_Handle, _Pid) ->
erlang:nif_error(nif_not_loaded).
nif_controlling_process_conn(_Handle, _Pid) ->
erlang:nif_error(nif_not_loaded).
%%%===================================================================
%%% Internal
%%%===================================================================
to_binary(B) when is_binary(B) -> B;
to_binary(L) when is_list(L) -> list_to_binary(L);
to_binary(A) when is_atom(A) -> atom_to_binary(A).