Packages
macula
4.2.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/peering/macula_peering_conn.erl
%% @doc Per-peer connection state machine.
%%
%% Implements the lifecycle from `Part 4 §10' simplified for Phase 1:
%% no REFRESH phase, no RECONNECTING — just enough to exchange signed
%% CONNECT/HELLO frames and drain on GOODBYE.
%%
%% State graph:
%% <pre>
%% client: connecting → handshaking → connected → draining → (terminate)
%% server: awaiting_start → handshaking → connected → draining → (terminate)
%% </pre>
-module(macula_peering_conn).
-behaviour(gen_statem).
-export([start_link/1]).
-export([init/1, callback_mode/0, terminate/3, code_change/4]).
-export([connecting/3, awaiting_start/3, handshaking/3, connected/3, draining/3]).
-export_type([opts/0, connect_opts/0]).
-type connect_opts() :: #{
host := binary() | string(),
port := inet:port_number(),
alpn => [binary()],
timeout_ms => timeout(),
_ => _
}.
-type opts() :: #{
role := client | server,
identity := macula_identity:key_pair(),
realms := [macula_identity:pubkey()],
capabilities := non_neg_integer(),
controlling_pid := pid(),
target => connect_opts(),
quic_conn => reference(),
%% Optional pid notified once when the worker completes the
%% CONNECT/HELLO handshake and transitions to `connected'. Sent
%% as `{macula_peering, handshake_complete, self(), PeerNodeId}'
%% where `PeerNodeId' is the verified peer Ed25519 pubkey from
%% the inbound CONNECT/HELLO frame. Used by accept-side listeners
%% that (a) cap concurrent *handshaking* workers and need to
%% release a slot the moment a worker is verified-and-connected,
%% and (b) dedupe duplicate dials from the same peer identity by
%% closing prior workers for the same `PeerNodeId'. Distinct from
%% `controlling_pid', which receives the peer-node-id-bearing
%% `connected' / `frame' / `disconnected' stream.
accept_owner => pid()
}.
-record(data, {
role :: client | server,
identity :: macula_identity:key_pair(),
node_id :: macula_identity:pubkey(),
realms :: [macula_identity:pubkey()],
capabilities :: non_neg_integer(),
controlling_pid :: pid(),
accept_owner :: undefined | pid(),
target :: undefined | connect_opts(),
quic_conn :: undefined | reference(),
quic_stream :: undefined | reference(),
peer_node_id :: undefined | macula_identity:pubkey(),
peer_station_id :: undefined | macula_identity:pubkey(),
peer_realms :: [macula_identity:pubkey()],
buf :: binary()
}).
-define(DRAIN_TIMEOUT_MS, 5_000).
%% Maximum time the `handshaking' state may take before the worker
%% gives up. CONNECT/HELLO is sub-second on a healthy peer; 30s is
%% generous. Drains workers stuck because the peer speaks the wrong
%% protocol (e.g. V1 frames against a V2 station) — without this the
%% sup accumulates stuck workers indefinitely. See PLAN_FLYING_RESTART.
-define(HANDSHAKE_TIMEOUT_MS, 30_000).
%%------------------------------------------------------------------
%% Lifecycle
%%------------------------------------------------------------------
-spec start_link(opts()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_statem:start_link(?MODULE, Opts, []).
callback_mode() ->
[state_functions, state_enter].
init(#{role := Role, identity := Identity, controlling_pid := Pid} = Opts)
when Role =:= client; Role =:= server ->
Data = #data{
role = Role,
identity = Identity,
node_id = macula_identity:public(Identity),
realms = maps:get(realms, Opts, []),
capabilities = maps:get(capabilities, Opts, 0),
controlling_pid = Pid,
accept_owner = maps:get(accept_owner, Opts, undefined),
target = maps:get(target, Opts, undefined),
quic_conn = maps:get(quic_conn, Opts, undefined),
quic_stream = undefined,
peer_node_id = undefined,
peer_station_id = undefined,
peer_realms = [],
buf = <<>>
},
{ok, initial_state(Role), Data}.
initial_state(client) -> connecting;
initial_state(server) -> awaiting_start.
terminate(_Reason, _State, Data) ->
_ = close_quic(Data),
ok.
code_change(_OldVsn, State, Data, _Extra) ->
{ok, State, Data}.
%%------------------------------------------------------------------
%% State: connecting (client only)
%%------------------------------------------------------------------
connecting(enter, _Old, Data) ->
self() ! attempt_connect,
{keep_state, Data};
connecting(info, attempt_connect, #data{target = Target} = Data) ->
after_connect(do_connect(Target), Data);
connecting(cast, {close, Reason}, Data) ->
notify(disconnected, Reason, Data),
{stop, normal, Data};
connecting(EventType, Event, Data) ->
drop_unexpected(EventType, Event, connecting, Data).
after_connect({ok, Conn}, Data) ->
ok = macula_quic:controlling_process(Conn, self()),
{next_state, handshaking, Data#data{quic_conn = Conn}};
after_connect(Other, Data) ->
notify(disconnected, {connect_failed, Other}, Data),
{stop, normal, Data}.
%%------------------------------------------------------------------
%% State: awaiting_start (server only — wait for ownership transfer)
%%------------------------------------------------------------------
awaiting_start(enter, _Old, Data) ->
{keep_state, Data};
awaiting_start(cast, start_handshake, Data) ->
{next_state, handshaking, Data};
awaiting_start(cast, {close, Reason}, Data) ->
notify(disconnected, Reason, Data),
{stop, normal, Data};
awaiting_start(EventType, Event, Data) ->
drop_unexpected(EventType, Event, awaiting_start, Data).
%%------------------------------------------------------------------
%% State: handshaking
%%------------------------------------------------------------------
handshaking(enter, _Old, #data{role = client, quic_conn = Conn} = Data) ->
on_handshake_enter_client(macula_quic:open_stream(Conn), Data);
handshaking(enter, _Old, #data{role = server, quic_conn = Conn} = Data) ->
ok = macula_quic:async_accept_stream(Conn),
{keep_state, Data, [handshake_state_timeout()]};
handshaking(info, {quic, new_stream, Stream, _Info}, Data) ->
ok = macula_quic:setopt(Stream, active, true),
{keep_state, Data#data{quic_stream = Stream}};
handshaking(info, {quic, Bin, Stream, _Flags},
#data{quic_stream = Stream, buf = Buf} = Data) when is_binary(Bin) ->
consume_handshake(<<Buf/binary, Bin/binary>>, Data);
handshaking(info, {quic, closed, _Conn, _Detail}, Data) ->
notify(disconnected, closed_during_handshake, Data),
{stop, normal, Data};
handshaking(cast, {close, Reason}, Data) ->
notify(disconnected, Reason, Data),
{stop, normal, Data};
%% No CONNECT/HELLO completed within the timeout window. Most common
%% cause: peer is speaking a different protocol version (e.g. V1
%% frames at a V2 station) — bytes accumulate in `buf' but never form
%% a valid frame. Surface a structured diagnostic and exit so the sup
%% does not retain the worker forever.
handshaking(state_timeout, handshake_timeout,
#data{role = Role, buf = Buf, quic_stream = Stream} = Data) ->
macula_diagnostics:event(<<"_macula.peering.handshake_timeout">>, #{
role => Role,
buf_size => byte_size(Buf),
has_stream => Stream =/= undefined,
timeout_ms => ?HANDSHAKE_TIMEOUT_MS
}),
notify(disconnected, handshake_timeout, Data),
{stop, normal, Data};
handshaking(EventType, Event, Data) ->
drop_unexpected(EventType, Event, handshaking, Data).
handshake_state_timeout() ->
{state_timeout, ?HANDSHAKE_TIMEOUT_MS, handshake_timeout}.
on_handshake_enter_client({ok, Stream}, Data) ->
%% setopt/send can both fail if the QUIC connection died between
%% nif_connect returning {ok, Conn} and us getting here (peer
%% closed, network drop, server rejected with a CONNECTION_CLOSE
%% frame after the TLS handshake but before we open a stream).
%% Prior to 3.15.3 the `ok = ...` matches turned every such
%% race into a crash; now we surface a structured disconnect
%% and let the caller schedule a reconnect.
case macula_quic:setopt(Stream, active, true) of
ok ->
case send_connect(Stream, Data) of
ok ->
{keep_state,
Data#data{quic_stream = Stream},
[handshake_state_timeout()]};
{error, _} = SendErr ->
notify(disconnected, {send_connect_failed, SendErr}, Data),
{stop, normal, Data}
end;
{error, _} = SetoptErr ->
notify(disconnected, {setopt_failed, SetoptErr}, Data),
{stop, normal, Data}
end;
on_handshake_enter_client(Err, Data) ->
notify(disconnected, {open_stream_failed, Err}, Data),
{stop, normal, Data}.
consume_handshake(Buf, Data) ->
{Frames, Tail} = macula_frame:parse_stream(Buf),
handle_handshake_frames(Frames, Data#data{buf = Tail}).
handle_handshake_frames([], Data) ->
{keep_state, Data};
handle_handshake_frames([#{frame_type := connect} = F | _], Data) ->
process_connect(F, Data);
handle_handshake_frames([#{frame_type := hello} = F | _], Data) ->
process_hello(F, Data);
handle_handshake_frames([_Other | Rest], Data) ->
handle_handshake_frames(Rest, Data).
%% Server side: peer's CONNECT
process_connect(#{node_id := PeerNodeId} = Frame,
#data{role = server, quic_stream = Stream} = Data) ->
on_connect_verified(macula_frame:verify(Frame, PeerNodeId), Frame, Stream, Data);
process_connect(_Frame, Data) ->
notify(disconnected, unexpected_connect_on_client, Data),
{stop, normal, Data}.
on_connect_verified({ok, _Verified}, Frame, Stream, Data) ->
NewData = absorb_peer_info(Frame, Data),
ok = send_hello(Stream, NewData),
transition_to_connected(NewData);
on_connect_verified({error, R}, _Frame, _Stream, Data) ->
notify(disconnected, {connect_verify_failed, R}, Data),
{stop, normal, Data}.
%% Client side: peer's HELLO
process_hello(#{node_id := PeerNodeId} = Frame, #data{role = client} = Data) ->
on_hello_verified(macula_frame:verify(Frame, PeerNodeId), Frame, Data);
process_hello(_Frame, Data) ->
notify(disconnected, unexpected_hello_on_server, Data),
{stop, normal, Data}.
on_hello_verified({ok, _Verified}, #{accepted := true} = Frame, Data) ->
transition_to_connected(absorb_peer_info(Frame, Data));
on_hello_verified({ok, _Verified}, #{accepted := false} = Frame, Data) ->
notify(disconnected, {refused, maps:get(refusal_code, Frame, undefined)}, Data),
{stop, normal, Data};
on_hello_verified({error, R}, _Frame, Data) ->
notify(disconnected, {hello_verify_failed, R}, Data),
{stop, normal, Data}.
absorb_peer_info(Frame, Data) ->
Data#data{
peer_node_id = maps:get(node_id, Frame),
peer_station_id = maps:get(station_id, Frame),
peer_realms = maps:get(realms, Frame, [])
}.
transition_to_connected(Data) ->
notify(connected, Data#data.peer_node_id, Data),
notify_handshake_complete(Data),
{next_state, connected, Data}.
notify_handshake_complete(#data{accept_owner = undefined}) ->
ok;
notify_handshake_complete(#data{accept_owner = Pid, peer_node_id = NodeId})
when is_pid(Pid) ->
Pid ! {macula_peering, handshake_complete, self(), NodeId},
ok.
%%------------------------------------------------------------------
%% State: connected
%%------------------------------------------------------------------
connected(enter, _Old, Data) ->
{keep_state, Data};
connected(info, {quic, Bin, Stream, _Flags},
#data{quic_stream = Stream, buf = Buf} = Data) when is_binary(Bin) ->
{Frames, Tail} = macula_frame:parse_stream(<<Buf/binary, Bin/binary>>),
[notify(frame, F, Data) || F <- Frames],
{keep_state, Data#data{buf = Tail}};
connected(info, {quic, closed, _Conn, _Detail}, Data) ->
notify(disconnected, peer_closed, Data),
{stop, normal, Data};
connected(cast, {close, Reason}, Data) ->
_ = send_goodbye(Data#data.quic_stream, Reason, Data),
{next_state, draining, Data};
connected(cast, {send_frame, Frame}, Data) ->
_ = send_application_frame(Frame, Data),
{keep_state, Data};
connected(EventType, Event, Data) ->
drop_unexpected(EventType, Event, connected, Data).
%%------------------------------------------------------------------
%% State: draining
%%------------------------------------------------------------------
draining(enter, _Old, Data) ->
{keep_state, Data, [{state_timeout, ?DRAIN_TIMEOUT_MS, drain_done}]};
draining(state_timeout, drain_done, Data) ->
notify(disconnected, drained, Data),
_ = close_quic(Data),
{stop, normal, Data};
draining(info, {quic, closed, _Conn, _Detail}, Data) ->
notify(disconnected, peer_closed_during_drain, Data),
{stop, normal, Data};
draining(info, {quic, _, _, _}, Data) ->
%% Ignore late inbound during drain.
{keep_state, Data};
draining(cast, {close, _Reason}, Data) ->
%% Already draining — idempotent.
{keep_state, Data};
draining(EventType, Event, Data) ->
drop_unexpected(EventType, Event, draining, Data).
%%------------------------------------------------------------------
%% Frame send helpers
%%------------------------------------------------------------------
send_connect(Stream, Data) ->
Frame = macula_frame:connect(#{
node_id => Data#data.node_id,
station_id => Data#data.node_id,
realms => Data#data.realms,
capabilities => Data#data.capabilities,
puzzle_evidence => macula_identity:puzzle_evidence(Data#data.node_id)
}),
Signed = macula_frame:sign(Frame, Data#data.identity),
macula_quic:send(Stream, macula_frame:encode(Signed)).
send_hello(Stream, Data) ->
Frame = macula_frame:hello(#{
node_id => Data#data.node_id,
station_id => Data#data.node_id,
realms => Data#data.realms,
capabilities => Data#data.capabilities,
accepted => true,
negotiated_capabilities => Data#data.capabilities
}),
Signed = macula_frame:sign(Frame, Data#data.identity),
macula_quic:send(Stream, macula_frame:encode(Signed)).
send_goodbye(undefined, _Reason, _Data) ->
ok;
send_goodbye(Stream, Reason, Data) ->
Frame = macula_frame:goodbye(Reason, undefined),
Signed = macula_frame:sign(Frame, Data#data.identity),
macula_quic:send(Stream, macula_frame:encode(Signed)).
send_application_frame(_Frame, #data{quic_stream = undefined}) ->
ok;
send_application_frame(Frame, #data{quic_stream = Stream, identity = Id}) ->
Signed = ensure_signed(Frame, Id),
macula_quic:send(Stream, macula_frame:encode(Signed)).
ensure_signed(#{signature := _} = Frame, _Id) -> Frame;
ensure_signed(Frame, Id) -> macula_frame:sign(Frame, Id).
close_quic(#data{quic_conn = undefined}) ->
ok;
close_quic(#data{quic_conn = Conn}) ->
catch macula_quic:close_connection(Conn),
ok.
%%------------------------------------------------------------------
%% Outbound dial — unpack peering's option-map into macula_quic's
%% positional API.
%%------------------------------------------------------------------
do_connect(#{host := Host, port := Port} = Target) ->
Timeout = maps:get(timeout_ms, Target, 30_000),
Alpn = maps:get(alpn, Target, [<<"macula">>]),
macula_quic:connect(Host, Port, [{alpn, Alpn}], Timeout).
%%------------------------------------------------------------------
%% Notifications
%%------------------------------------------------------------------
notify(Event, Detail, #data{controlling_pid = Pid}) ->
Pid ! {macula_peering, Event, self(), Detail},
ok.
drop_unexpected(EventType, Event, State, Data) ->
macula_diagnostics:event(<<"_macula.peering.unexpected">>, #{
state => State,
event_type => EventType,
event => safe_event(Event)
}),
{keep_state, Data}.
%% Truncate large/binary events for safer log emission.
safe_event(Bin) when is_binary(Bin), byte_size(Bin) > 64 ->
{truncated, byte_size(Bin)};
safe_event(Other) ->
Other.